From a1a70dea5d33cf5eb595fdce3dac1c3ed963fed4 Mon Sep 17 00:00:00 2001 From: saopig1 <4x7sw862st@gmail.com> Date: Thu, 25 Jun 2026 16:10:38 +0800 Subject: [PATCH] fix(guest): serialize concurrent queue-mutation playback per bot The queue-mutating playback routes (play-now-song, play-next-song, add-song, play-at) read queue position synchronously, mutate the queue, then await resolveAndPlay() which suspends at an async URL fetch before player.play(). With no serialization, two concurrent requests (normal in login-less guest mode) interleave: the audible song (decided by URL-fetch latency) can disagree with queue.currentIndex (decided by sync-block ordering), corrupting "now playing" and causing skipped/duplicate songs. Add a per-bot async serializer (BotInstance.runExclusive) and wrap the critical region of all four routes in it. Single-request behavior and every response shape / validation 400 are preserved; only the critical region moved inside runExclusive. Co-Authored-By: Claude Opus 4.8 (1M context) --- src/bot/instance.test.ts | 112 +++++++++++++++++++++++++++++ src/bot/instance.ts | 9 +++ src/web/api/player.ts | 148 ++++++++++++++++++++++----------------- 3 files changed, 204 insertions(+), 65 deletions(-) create mode 100644 src/bot/instance.test.ts diff --git a/src/bot/instance.test.ts b/src/bot/instance.test.ts new file mode 100644 index 0000000..751116c --- /dev/null +++ b/src/bot/instance.test.ts @@ -0,0 +1,112 @@ +import { describe, it, expect } from "vitest"; +import { BotInstance } from "./instance.js"; + +// Constructing a real BotInstance is heavy (spawns a TS3Client, AudioPlayer, +// reads avatars, etc.), and runExclusive only touches a single private field +// (`playGate`). So we exercise the ACTUAL shipped method via its prototype, +// bound to a minimal object carrying just that field. This proves the real +// serializer logic without standing up a full bot. +type Gate = { playGate: Promise }; +const runExclusive = BotInstance.prototype.runExclusive as ( + this: Gate, + fn: () => Promise, +) => Promise; + +function makeGate(): Gate { + return { playGate: Promise.resolve() }; +} + +/** An explicit, timer-free deferred so ordering is deterministic. */ +function deferred() { + let resolve!: (value: T) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + return { promise, resolve, reject }; +} + +describe("BotInstance.runExclusive — serialization", () => { + it("does not start fnB until fnA settles", async () => { + const gate = makeGate(); + const order: string[] = []; + const gateA = deferred(); + + const pA = runExclusive.call(gate, async () => { + order.push("A-start"); + await gateA.promise; // suspend A until we explicitly release it + order.push("A-end"); + }); + + const pB = runExclusive.call(gate, async () => { + order.push("B-start"); + order.push("B-end"); + }); + + // Give the microtask queue a chance: B must NOT have started while A is + // still suspended on gateA. + await Promise.resolve(); + await Promise.resolve(); + expect(order).toEqual(["A-start"]); + + gateA.resolve(); + await pA; + await pB; + + expect(order).toEqual(["A-start", "A-end", "B-start", "B-end"]); + }); + + it("runs fnB even if fnA rejects (chain survives rejection)", async () => { + const gate = makeGate(); + const order: string[] = []; + const gateA = deferred(); + + const pA = runExclusive.call(gate, async () => { + order.push("A-start"); + await gateA.promise; + throw new Error("A blew up"); + }); + + const pB = runExclusive.call(gate, async () => { + order.push("B-start"); + order.push("B-end"); + return "B-result"; + }); + + await Promise.resolve(); + await Promise.resolve(); + expect(order).toEqual(["A-start"]); + + gateA.reject(new Error("A blew up")); + await expect(pA).rejects.toThrow("A blew up"); + + // B still runs, only after A has fully settled. + await expect(pB).resolves.toBe("B-result"); + expect(order).toEqual(["A-start", "B-start", "B-end"]); + }); + + it("preserves call order across three serialized tasks", async () => { + const gate = makeGate(); + const order: string[] = []; + const tasks = ["X", "Y", "Z"]; + const promises = tasks.map((t) => + runExclusive.call(gate, async () => { + order.push(`${t}-start`); + await Promise.resolve(); + order.push(`${t}-end`); + }), + ); + + await Promise.all(promises); + + expect(order).toEqual([ + "X-start", + "X-end", + "Y-start", + "Y-end", + "Z-start", + "Z-end", + ]); + }); +}); diff --git a/src/bot/instance.ts b/src/bot/instance.ts index 0e19128..e83ec25 100755 --- a/src/bot/instance.ts +++ b/src/bot/instance.ts @@ -78,6 +78,7 @@ export class BotInstance extends EventEmitter { private fmProvider: MusicProvider | null = null; /** Results of the most recent !search, for "#N" selection (issue #90). */ private lastSearchResults: Song[] = []; + private playGate: Promise = Promise.resolve(); constructor(options: BotInstanceOptions) { super(); @@ -1051,6 +1052,14 @@ export class BotInstance extends EventEmitter { return input; } + /** Serialize queue-mutation + play sequences so concurrent requests can't + * interleave (audible track must match queue.currentIndex). */ + runExclusive(fn: () => Promise): Promise { + const next = this.playGate.then(fn, fn); + this.playGate = next.catch(() => {}); + return next; + } + getStatus(): BotStatus { return { id: this.id, diff --git a/src/web/api/player.ts b/src/web/api/player.ts index 4811d8f..23fd5fe 100644 --- a/src/web/api/player.ts +++ b/src/web/api/player.ts @@ -215,27 +215,34 @@ export function createPlayerRouter( res.status(400).json({ error: "index is required" }); return; } - const queue = bot.getQueueManager(); - // Validate the index BEFORE stopping current playback — otherwise an - // invalid index silently kills the user's current song and leaves the - // queue idle. - if (index >= queue.size()) { - res.status(400).json({ error: "Invalid queue index" }); + // Serialize the index-validation + stop/reset/playAt/resolveAndPlay so a + // concurrent request can't interleave between mutating the queue and + // starting playback (audible track must match queue.currentIndex). + const result = await bot.runExclusive(async () => { + const queue = bot.getQueueManager(); + // Validate the index BEFORE stopping current playback — otherwise an + // invalid index silently kills the user's current song and leaves the + // queue idle. + if (index >= queue.size()) { + return { status: 400 as const, body: { error: "Invalid queue index" } }; + } + bot.getPlayer().stop(); + bot.getPlayer().resetFailures(); + const song = queue.playAt(index); + if (!song) { + return { status: 400 as const, body: { error: "Invalid queue index" } }; + } + const ok = await bot.resolveAndPlay(song); + if (!ok) { + return { body: { message: `Cannot play: ${song.name}` } }; + } + return { body: { message: `Now playing: ${song.name} - ${song.artist}` } }; + }); + if (result.status) { + res.status(result.status).json(result.body); return; } - bot.getPlayer().stop(); - bot.getPlayer().resetFailures(); - const song = queue.playAt(index); - if (!song) { - res.status(400).json({ error: "Invalid queue index" }); - return; - } - const ok = await bot.resolveAndPlay(song); - if (!ok) { - res.json({ message: `Cannot play: ${song.name}` }); - return; - } - res.json({ message: `Now playing: ${song.name} - ${song.artist}` }); + res.json(result.body); } catch (err) { res.status(500).json({ error: (err as Error).message }); } @@ -453,31 +460,34 @@ export function createPlayerRouter( res.status(400).json({ error: "song object with id and platform is required" }); return; } - const queue = bot.getQueueManager(); - const wasIdle = bot.getPlayer().getState() === "idle"; - // Capture the slot addNext WILL insert at, before mutating the queue. - // addNext pushes when currentIndex<0 (slot = size); otherwise splices - // at currentIndex+1. Using size-1 after addNext was wrong when the - // queue had stale currentIndex>=0 while the player was idle (e.g., - // after natural track end without queue.clear()). - const insertedAt = - queue.getCurrentIndex() < 0 ? queue.size() : queue.getCurrentIndex() + 1; - queue.addNext(song); + // Serialize the queue mutation + playback so concurrent requests can't + // interleave (audible track must match queue.currentIndex). + const body = await bot.runExclusive(async () => { + const queue = bot.getQueueManager(); + const wasIdle = bot.getPlayer().getState() === "idle"; + // Capture the slot addNext WILL insert at, before mutating the queue. + // addNext pushes when currentIndex<0 (slot = size); otherwise splices + // at currentIndex+1. Using size-1 after addNext was wrong when the + // queue had stale currentIndex>=0 while the player was idle (e.g., + // after natural track end without queue.clear()). + const insertedAt = + queue.getCurrentIndex() < 0 ? queue.size() : queue.getCurrentIndex() + 1; + queue.addNext(song); - if (wasIdle) { - // Promote the just-added song to current and start it. - queue.playAt(insertedAt); - bot.getPlayer().resetFailures(); - const ok = await bot.resolveAndPlay(queue.current()!); - if (!ok) { - res.json({ ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` }); - return; + if (wasIdle) { + // Promote the just-added song to current and start it. + queue.playAt(insertedAt); + bot.getPlayer().resetFailures(); + const ok = await bot.resolveAndPlay(queue.current()!); + if (!ok) { + return { ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` }; + } + return { ok: true, message: `正在播放:${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` }; } - res.json({ ok: true, message: `正在播放:${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` }); - return; - } - res.json({ ok: true, message: `已加入下一首:${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` }); + return { ok: true, message: `已加入下一首:${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` }; + }); + res.json(body); } catch (err) { res.status(500).json({ error: (err as Error).message }); } @@ -494,18 +504,22 @@ export function createPlayerRouter( res.status(400).json({ error: "song object with id and platform is required" }); return; } - const queue = bot.getQueueManager(); - const insertedAt = - queue.getCurrentIndex() < 0 ? queue.size() : queue.getCurrentIndex() + 1; - queue.addNext(song); - queue.playAt(insertedAt); - bot.getPlayer().resetFailures(); - const ok = await bot.resolveAndPlay(queue.current()!); - if (!ok) { - res.json({ ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` }); - return; - } - res.json({ ok: true, message: `正在播放:${song.name || "Unknown"} - ${song.artist || "Unknown"}` }); + // Serialize the insert-after-current + promote + playback so concurrent + // requests can't interleave (audible track must match queue.currentIndex). + const body = await bot.runExclusive(async () => { + const queue = bot.getQueueManager(); + const insertedAt = + queue.getCurrentIndex() < 0 ? queue.size() : queue.getCurrentIndex() + 1; + queue.addNext(song); + queue.playAt(insertedAt); + bot.getPlayer().resetFailures(); + const ok = await bot.resolveAndPlay(queue.current()!); + if (!ok) { + return { ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` }; + } + return { ok: true, message: `正在播放:${song.name || "Unknown"} - ${song.artist || "Unknown"}` }; + }); + res.json(body); } catch (err) { res.status(500).json({ error: (err as Error).message }); } @@ -519,20 +533,24 @@ export function createPlayerRouter( res.status(400).json({ error: "song object with id and platform is required" }); return; } - const queue = bot.getQueueManager(); - const wasIdle = bot.getPlayer().getState() === "idle"; - queue.add(song); + // Serialize the queue mutation + (possible) playback so concurrent + // requests can't interleave (audible track must match queue.currentIndex). + const body = await bot.runExclusive(async () => { + const queue = bot.getQueueManager(); + const wasIdle = bot.getPlayer().getState() === "idle"; + queue.add(song); - // If nothing was playing, start this newly-added song immediately. - if (wasIdle) { - queue.playAt(queue.size() - 1); - bot.getPlayer().resetFailures(); - await bot.resolveAndPlay(queue.current()!); - res.json({ message: `Now playing: ${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` }); - return; - } + // If nothing was playing, start this newly-added song immediately. + if (wasIdle) { + queue.playAt(queue.size() - 1); + bot.getPlayer().resetFailures(); + await bot.resolveAndPlay(queue.current()!); + return { message: `Now playing: ${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` }; + } - res.json({ message: `Added to queue: ${song.name || 'Unknown'} - ${song.artist || 'Unknown'} (position ${queue.size()})` }); + return { message: `Added to queue: ${song.name || 'Unknown'} - ${song.artist || 'Unknown'} (position ${queue.size()})` }; + }); + res.json(body); } catch (err) { res.status(500).json({ error: (err as Error).message }); }