mirror of
https://github.com/ZHANGTIANYAO1/teamspeak-music-bot.git
synced 2026-10-02 04:52:50 +08:00
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) <noreply@anthropic.com>
This commit is contained in:
1 parent
43c0175334
commit
a1a70dea5d
3 files changed
+156
-17
No files matched your search
@@ -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<unknown> };
|
||||||
|
const runExclusive = BotInstance.prototype.runExclusive as <T>(
|
||||||
|
this: Gate,
|
||||||
|
fn: () => Promise<T>,
|
||||||
|
) => Promise<T>;
|
||||||
|
|
||||||
|
function makeGate(): Gate {
|
||||||
|
return { playGate: Promise.resolve() };
|
||||||
|
}
|
||||||
|
|
||||||
|
/** An explicit, timer-free deferred so ordering is deterministic. */
|
||||||
|
function deferred<T = void>() {
|
||||||
|
let resolve!: (value: T) => void;
|
||||||
|
let reject!: (reason?: unknown) => void;
|
||||||
|
const promise = new Promise<T>((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",
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -78,6 +78,7 @@ export class BotInstance extends EventEmitter {
|
|||||||
private fmProvider: MusicProvider | null = null;
|
private fmProvider: MusicProvider | null = null;
|
||||||
/** Results of the most recent !search, for "#N" selection (issue #90). */
|
/** Results of the most recent !search, for "#N" selection (issue #90). */
|
||||||
private lastSearchResults: Song[] = [];
|
private lastSearchResults: Song[] = [];
|
||||||
|
private playGate: Promise<unknown> = Promise.resolve();
|
||||||
|
|
||||||
constructor(options: BotInstanceOptions) {
|
constructor(options: BotInstanceOptions) {
|
||||||
super();
|
super();
|
||||||
@@ -1051,6 +1052,14 @@ export class BotInstance extends EventEmitter {
|
|||||||
return input;
|
return input;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Serialize queue-mutation + play sequences so concurrent requests can't
|
||||||
|
* interleave (audible track must match queue.currentIndex). */
|
||||||
|
runExclusive<T>(fn: () => Promise<T>): Promise<T> {
|
||||||
|
const next = this.playGate.then(fn, fn);
|
||||||
|
this.playGate = next.catch(() => {});
|
||||||
|
return next;
|
||||||
|
}
|
||||||
|
|
||||||
getStatus(): BotStatus {
|
getStatus(): BotStatus {
|
||||||
return {
|
return {
|
||||||
id: this.id,
|
id: this.id,
|
||||||
|
|||||||
+35
-17
@@ -215,27 +215,34 @@ export function createPlayerRouter(
|
|||||||
res.status(400).json({ error: "index is required" });
|
res.status(400).json({ error: "index is required" });
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
// 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();
|
const queue = bot.getQueueManager();
|
||||||
// Validate the index BEFORE stopping current playback — otherwise an
|
// Validate the index BEFORE stopping current playback — otherwise an
|
||||||
// invalid index silently kills the user's current song and leaves the
|
// invalid index silently kills the user's current song and leaves the
|
||||||
// queue idle.
|
// queue idle.
|
||||||
if (index >= queue.size()) {
|
if (index >= queue.size()) {
|
||||||
res.status(400).json({ error: "Invalid queue index" });
|
return { status: 400 as const, body: { error: "Invalid queue index" } };
|
||||||
return;
|
|
||||||
}
|
}
|
||||||
bot.getPlayer().stop();
|
bot.getPlayer().stop();
|
||||||
bot.getPlayer().resetFailures();
|
bot.getPlayer().resetFailures();
|
||||||
const song = queue.playAt(index);
|
const song = queue.playAt(index);
|
||||||
if (!song) {
|
if (!song) {
|
||||||
res.status(400).json({ error: "Invalid queue index" });
|
return { status: 400 as const, body: { error: "Invalid queue index" } };
|
||||||
return;
|
|
||||||
}
|
}
|
||||||
const ok = await bot.resolveAndPlay(song);
|
const ok = await bot.resolveAndPlay(song);
|
||||||
if (!ok) {
|
if (!ok) {
|
||||||
res.json({ message: `Cannot play: ${song.name}` });
|
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;
|
return;
|
||||||
}
|
}
|
||||||
res.json({ message: `Now playing: ${song.name} - ${song.artist}` });
|
res.json(result.body);
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
res.status(500).json({ error: (err as Error).message });
|
res.status(500).json({ error: (err as Error).message });
|
||||||
}
|
}
|
||||||
@@ -453,6 +460,9 @@ export function createPlayerRouter(
|
|||||||
res.status(400).json({ error: "song object with id and platform is required" });
|
res.status(400).json({ error: "song object with id and platform is required" });
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
// 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 queue = bot.getQueueManager();
|
||||||
const wasIdle = bot.getPlayer().getState() === "idle";
|
const wasIdle = bot.getPlayer().getState() === "idle";
|
||||||
// Capture the slot addNext WILL insert at, before mutating the queue.
|
// Capture the slot addNext WILL insert at, before mutating the queue.
|
||||||
@@ -470,14 +480,14 @@ export function createPlayerRouter(
|
|||||||
bot.getPlayer().resetFailures();
|
bot.getPlayer().resetFailures();
|
||||||
const ok = await bot.resolveAndPlay(queue.current()!);
|
const ok = await bot.resolveAndPlay(queue.current()!);
|
||||||
if (!ok) {
|
if (!ok) {
|
||||||
res.json({ ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` });
|
return { ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` };
|
||||||
return;
|
|
||||||
}
|
}
|
||||||
res.json({ ok: true, message: `正在播放:${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` });
|
return { 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) {
|
} catch (err) {
|
||||||
res.status(500).json({ error: (err as Error).message });
|
res.status(500).json({ error: (err as Error).message });
|
||||||
}
|
}
|
||||||
@@ -494,6 +504,9 @@ export function createPlayerRouter(
|
|||||||
res.status(400).json({ error: "song object with id and platform is required" });
|
res.status(400).json({ error: "song object with id and platform is required" });
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
// 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 queue = bot.getQueueManager();
|
||||||
const insertedAt =
|
const insertedAt =
|
||||||
queue.getCurrentIndex() < 0 ? queue.size() : queue.getCurrentIndex() + 1;
|
queue.getCurrentIndex() < 0 ? queue.size() : queue.getCurrentIndex() + 1;
|
||||||
@@ -502,10 +515,11 @@ export function createPlayerRouter(
|
|||||||
bot.getPlayer().resetFailures();
|
bot.getPlayer().resetFailures();
|
||||||
const ok = await bot.resolveAndPlay(queue.current()!);
|
const ok = await bot.resolveAndPlay(queue.current()!);
|
||||||
if (!ok) {
|
if (!ok) {
|
||||||
res.json({ ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` });
|
return { ok: false, message: `无法播放「${song.name || song.id}」(区域/版权限制)` };
|
||||||
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) {
|
} catch (err) {
|
||||||
res.status(500).json({ error: (err as Error).message });
|
res.status(500).json({ error: (err as Error).message });
|
||||||
}
|
}
|
||||||
@@ -519,6 +533,9 @@ export function createPlayerRouter(
|
|||||||
res.status(400).json({ error: "song object with id and platform is required" });
|
res.status(400).json({ error: "song object with id and platform is required" });
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
// 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 queue = bot.getQueueManager();
|
||||||
const wasIdle = bot.getPlayer().getState() === "idle";
|
const wasIdle = bot.getPlayer().getState() === "idle";
|
||||||
queue.add(song);
|
queue.add(song);
|
||||||
@@ -528,11 +545,12 @@ export function createPlayerRouter(
|
|||||||
queue.playAt(queue.size() - 1);
|
queue.playAt(queue.size() - 1);
|
||||||
bot.getPlayer().resetFailures();
|
bot.getPlayer().resetFailures();
|
||||||
await bot.resolveAndPlay(queue.current()!);
|
await bot.resolveAndPlay(queue.current()!);
|
||||||
res.json({ message: `Now playing: ${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` });
|
return { message: `Now playing: ${song.name || 'Unknown'} - ${song.artist || 'Unknown'}` };
|
||||||
return;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
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) {
|
} catch (err) {
|
||||||
res.status(500).json({ error: (err as Error).message });
|
res.status(500).json({ error: (err as Error).message });
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in new issue
Block a user