diff --git a/README.md b/README.md index 0e4e4a5..6544b48 100644 --- a/README.md +++ b/README.md @@ -953,7 +953,18 @@ 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.1:长视频播放阻塞与音乐 API 日志隐私修复 +### 最新版本 — 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 下载失败日志采用同样的处理。 diff --git a/src/audio/player.test.ts b/src/audio/player.test.ts index 9dd96ca..74531f3 100644 --- a/src/audio/player.test.ts +++ b/src/audio/player.test.ts @@ -40,6 +40,44 @@ 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); @@ -825,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; @@ -837,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; @@ -846,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 { diff --git a/src/audio/player.ts b/src/audio/player.ts index c526062..6892d6e 100644 --- a/src/audio/player.ts +++ b/src/audio/player.ts @@ -87,8 +87,20 @@ export function cleanupTempDir(dir: string): void { export function buildFfmpegArgs(url: string, seekSeconds: number): string[] { const args: string[] = ["-nostats"]; - const isHttp = /^https?:\/\//i.test(url); - const isBilibili = isHttp && (url.includes("bilivideo") || url.includes("bilibili")); + 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( @@ -662,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仍在运行但缓冲区不足一帧,且连续多次无法获取数据 @@ -683,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 diff --git a/src/bot/instance-recovery.test.ts b/src/bot/instance-recovery.test.ts new file mode 100644 index 0000000..6f57aec --- /dev/null +++ b/src/bot/instance-recovery.test.ts @@ -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() { + let resolve!: (value: T) => void; + const promise = new Promise(done => { resolve = done; }); + return { promise, resolve }; +} + +async function flush() { + await new Promise(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(), newer = deferred(); + 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(); + 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(), newer = deferred(); + 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(); + 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(); + 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(); + 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(); + }); +}); diff --git a/src/bot/instance.test.ts b/src/bot/instance.test.ts index d55292e..685a1e7 100644 --- a/src/bot/instance.test.ts +++ b/src/bot/instance.test.ts @@ -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; @@ -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([]); diff --git a/src/bot/instance.ts b/src/bot/instance.ts index cd323ff..5b6fe03 100755 --- a/src/bot/instance.ts +++ b/src/bot/instance.ts @@ -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(); 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 = 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 { 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 { + 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>; + 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(resolve => setTimeout(resolve, lookup * 1000)); + if (!isCurrent()) return cancelled(); + } + let result: Awaited> = 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 { @@ -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(), diff --git a/src/bot/voice-failure.test.ts b/src/bot/voice-failure.test.ts new file mode 100644 index 0000000..a81bd9f --- /dev/null +++ b/src/bot/voice-failure.test.ts @@ -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 }).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(); } + }); +}); diff --git a/src/ts-protocol/client-voice.test.ts b/src/ts-protocol/client-voice.test.ts new file mode 100644 index 0000000..feb0122 --- /dev/null +++ b/src/ts-protocol/client-voice.test.ts @@ -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(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[]) => 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(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"); + }); +}); diff --git a/src/ts-protocol/client.ts b/src/ts-protocol/client.ts index aa8a4c5..aded62a 100644 --- a/src/ts-protocol/client.ts +++ b/src/ts-protocol/client.ts @@ -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 | null = null; - private readonly voiceEndpointResolver = new TrackingVoiceEndpointResolver(); + private connectionGeneration = 0; + private voiceFailureTimer: ReturnType | 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 { - 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 { - 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): 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"); } }