fix: preserve pause intent and deduplicate stream recovery lookups

This commit is contained in:
TIANYAO ZHANG committed 2026-10-03 17:17:53 +08:00
1 parent cf83916f39
commit 9bcddb1881
2 files changed
+138 -15

No files matched your search

+105 -10
View File
@@ -3,6 +3,7 @@ import { EventEmitter } from "node:events";
import { BotInstance, COMMAND_DENIED_MESSAGE, spotifyPortsForBotId } from "./instance.js"; import { BotInstance, COMMAND_DENIED_MESSAGE, spotifyPortsForBotId } from "./instance.js";
import type { BotInstanceOptions } from "./instance.js"; import type { BotInstanceOptions } from "./instance.js";
import { PlayQueue, PlayMode } from "../audio/queue.js"; import { PlayQueue, PlayMode } from "../audio/queue.js";
import { AudioPlayer } from "../audio/player.js";
import { createDatabase, SHARED_QUEUE_OWNER } from "../data/database.js"; import { createDatabase, SHARED_QUEUE_OWNER } from "../data/database.js";
import { parseCommand } from "./commands.js"; import { parseCommand } from "./commands.js";
import type { TS3TextMessage } from "../ts-protocol/client.js"; import type { TS3TextMessage } from "../ts-protocol/client.js";
@@ -1760,13 +1761,15 @@ describe("BotInstance trackEnd — stale playback sessions", () => {
platform, duration, url: "old", platform, duration, url: "old",
}; };
let current: any = song; let current: any = song;
let state = "idle";
let session = 1;
const player = new EventEmitter() as any; const player = new EventEmitter() as any;
player.getState = () => state; player.state = "idle";
player.sessionId = 1;
player.getState = AudioPlayer.prototype.getState;
player.getElapsed = () => 1000; player.getElapsed = () => 1000;
player.getPlaybackSessionId = () => session; player.getPlaybackSessionId = AudioPlayer.prototype.getPlaybackSessionId;
player.play = vi.fn(() => { session++; state = "playing"; }); player.pause = AudioPlayer.prototype.pause;
player.resume = AudioPlayer.prototype.resume;
player.play = vi.fn(() => { player.sessionId++; player.state = "playing"; });
const provider = { getSongUrl: vi.fn(async () => ({ url: "fresh" })) }; const provider = { getSongUrl: vi.fn(async () => ({ url: "fresh" })) };
const advances: string[] = []; const advances: string[] = [];
const ctx: any = { const ctx: any = {
@@ -1776,10 +1779,11 @@ describe("BotInstance trackEnd — stale playback sessions", () => {
logger: { warn: vi.fn(), debug: vi.fn(), error: vi.fn() }, emit: vi.fn(), logger: { warn: vi.fn(), debug: vi.fn(), error: vi.fn() }, emit: vi.fn(),
getProviderFor: () => provider, getProviderFor: () => provider,
playNext: vi.fn(async () => { advances.push(current?.id ?? "empty"); return true; }), playNext: vi.fn(async () => { advances.push(current?.id ?? "empty"); return true; }),
replace: () => { current = { ...song, id: "replacement" }; session++; state = "playing"; }, replace: () => { current = { ...song, id: "replacement" }; player.sessionId++; player.state = "playing"; },
stop: () => { current = null; session++; state = "idle"; }, stop: () => { current = null; player.sessionId++; player.state = "idle"; },
restartSameSong: () => { session++; state = "idle"; }, restartSameSong: () => { player.sessionId++; player.state = "idle"; },
pause: () => { state = "paused"; }, pause: () => cmdPause.call(ctx),
resume: () => cmdResume.call(ctx),
advances, advances,
}; };
ctx.resumeInterruptedStream = (BotInstance.prototype as any).resumeInterruptedStream.bind(ctx); ctx.resumeInterruptedStream = (BotInstance.prototype as any).resumeInterruptedStream.bind(ctx);
@@ -1818,7 +1822,7 @@ describe("BotInstance trackEnd — stale playback sessions", () => {
expect(ctx.advances).toEqual([]); expect(ctx.advances).toEqual([]);
}); });
it.each(["stop", "restartSameSong", "pause"])("recovery does not overwrite playback after %s", async action => { it.each(["stop", "restartSameSong"])("recovery does not overwrite playback after %s", async action => {
const ctx = makeEndedCtx(); const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>(); const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValue(lookup.promise); ctx.provider.getSongUrl.mockReturnValue(lookup.promise);
@@ -1830,6 +1834,97 @@ describe("BotInstance trackEnd — stale playback sessions", () => {
expect(ctx.advances).toEqual([]); expect(ctx.advances).toEqual([]);
}); });
it("pause during an idle URL lookup is honored by recovered playback, then resume continues", async () => {
const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValue(lookup.promise);
ctx.player.emit("trackEnd");
ctx.pause();
expect(ctx.player.getState()).toBe("idle"); // actual AudioPlayer.pause cannot pause idle
lookup.resolve({ url: "fresh" });
await flushEvents();
expect(ctx.player.play).toHaveBeenCalledWith("fresh", 1000, 10_000);
expect(ctx.player.getState()).toBe("paused");
expect(ctx.advances).toEqual([]);
ctx.resume();
expect(ctx.player.getState()).toBe("playing");
expect(ctx.player.play).toHaveBeenCalledTimes(1);
});
it("resume before a paused recovery lookup completes lets the fresh stream play", async () => {
const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValue(lookup.promise);
ctx.player.emit("trackEnd");
ctx.pause();
ctx.resume();
lookup.resolve({ url: "fresh" });
await flushEvents();
expect(ctx.player.getState()).toBe("playing");
expect(ctx.advances).toEqual([]);
expect(ctx.provider.getSongUrl).toHaveBeenCalledTimes(1);
expect(ctx.streamRecovery?.attempts).toBe(1);
});
it("resume during a pending lookup cannot start a failing duplicate and skip the song", async () => {
const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValueOnce(lookup.promise).mockResolvedValue(null);
ctx.player.emit("trackEnd");
ctx.pause();
ctx.resume();
await flushEvents();
expect(ctx.advances).toEqual([]);
expect(ctx.provider.getSongUrl).toHaveBeenCalledTimes(1);
lookup.resolve({ url: "fresh" });
await flushEvents();
expect(ctx.player.getState()).toBe("playing");
expect(ctx.streamRecovery?.attempts).toBe(1);
});
it("pause while a recovery lookup fails prevents automatic queue advancement", async () => {
const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValue(lookup.promise);
ctx.player.emit("trackEnd");
ctx.pause();
lookup.reject(new Error("temporary lookup failure"));
await flushEvents();
expect(ctx.advances).toEqual([]);
});
it("resume after a paused failed lookup retries recovery instead of remaining idle", async () => {
const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValue(lookup.promise);
ctx.player.emit("trackEnd");
ctx.pause();
lookup.reject(new Error("temporary lookup failure"));
await flushEvents();
ctx.provider.getSongUrl.mockResolvedValue({ url: "recovered" });
ctx.resume();
await flushEvents();
expect(ctx.player.getState()).toBe("playing");
expect(ctx.player.play).toHaveBeenCalledWith("recovered", 1000, 10_000);
expect(ctx.advances).toEqual([]);
});
it("a same-song restart cannot inherit pause intent from an older rejected lookup", async () => {
const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValue(lookup.promise);
ctx.player.emit("trackEnd");
ctx.pause();
ctx.restartSameSong();
lookup.reject(new Error("temporary lookup failure"));
await flushEvents();
ctx.provider.getSongUrl.mockResolvedValue({ url: "new-recovery" });
ctx.player.emit("trackEnd");
await flushEvents();
expect(ctx.player.getState()).toBe("playing");
expect(ctx.advances).toEqual([]);
});
it("a failed recovery cannot advance a newer session of the same queue song", async () => { it("a failed recovery cannot advance a newer session of the same queue song", async () => {
const ctx = makeEndedCtx(); const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>(); const lookup = deferred<{ url: string }>();
+33 -5
View File
@@ -200,7 +200,7 @@ export class BotInstance extends EventEmitter {
/** 当前曲实际播放时长(试听片段秒数或完整 duration);resolveAndPlay 赋值。 */ /** 当前曲实际播放时长(试听片段秒数或完整 duration);resolveAndPlay 赋值。 */
private effectiveDuration: number | undefined; private effectiveDuration: number | undefined;
/** Resume attempts for the current song's stream (#161); see resumeInterruptedStream. */ /** Resume attempts for the current song's stream (#161); see resumeInterruptedStream. */
private streamRecovery: { song: QueuedSong; attempts: number; position: number } | null = null; private streamRecovery: { song: QueuedSong; attempts: number; position: number; session: number; pauseRequested: boolean; inFlight: boolean } | null = null;
private playGate: Promise<unknown> = Promise.resolve(); private playGate: Promise<unknown> = Promise.resolve();
/** Per-bot Jellyfin playback-report session (start / ~10s progress / stop). /** Per-bot Jellyfin playback-report session (start / ~10s progress / stop).
* null when the wired provider has no reporting capability. */ * null when the wired provider has no reporting capability. */
@@ -345,7 +345,8 @@ export class BotInstance extends EventEmitter {
!this.connected || !this.connected ||
this.queue.current() !== endedSong || this.queue.current() !== endedSong ||
this.player.getPlaybackSessionId() !== endedSession || this.player.getPlaybackSessionId() !== endedSession ||
this.player.getState() !== "idle" this.player.getState() !== "idle" ||
(this.streamRecovery?.song === endedSong && this.streamRecovery.pauseRequested)
) return; ) return;
this.logger.debug("Track ended, advancing queue"); this.logger.debug("Track ended, advancing queue");
return this.playNext(); return this.playNext();
@@ -1180,11 +1181,13 @@ export class BotInstance extends EventEmitter {
if ( if (
!recovery || !recovery ||
recovery.song !== song || recovery.song !== song ||
recovery.session !== endedSession ||
position - recovery.position > BotInstance.STREAM_END_TOLERANCE_S position - recovery.position > BotInstance.STREAM_END_TOLERANCE_S
) { ) {
this.streamRecovery = { song, attempts: 0, position }; this.streamRecovery = { song, attempts: 0, position, session: endedSession, pauseRequested: false, inFlight: false };
} }
const state = this.streamRecovery!; const state = this.streamRecovery!;
if (state.inFlight) return true;
if (state.attempts >= BotInstance.MAX_STREAM_RESUMES) { if (state.attempts >= BotInstance.MAX_STREAM_RESUMES) {
this.logger.warn( this.logger.warn(
{ songId: song.id, position, duration, attempts: state.attempts }, { songId: song.id, position, duration, attempts: state.attempts },
@@ -1200,18 +1203,31 @@ export class BotInstance extends EventEmitter {
{ songId: song.id, position, duration, attempt: state.attempts }, { songId: song.id, position, duration, attempt: state.attempts },
"Stream ended before the track did — resuming with a fresh URL", "Stream ended before the track did — resuming with a fresh URL",
); );
const result = await this.getProviderFor(song.platform).getSongUrl(song.id); state.inFlight = true;
let result: Awaited<ReturnType<MusicProvider["getSongUrl"]>>;
try {
result = await this.getProviderFor(song.platform).getSongUrl(song.id);
} finally {
state.inFlight = false;
}
// The user may have skipped/stopped while we were resolving; never // The user may have skipped/stopped while we were resolving; never
// clobber whatever is playing now. // clobber whatever is playing now.
if ( if (
this.queue.current() !== song || this.queue.current() !== song ||
this.player.getPlaybackSessionId() !== endedSession || this.player.getPlaybackSessionId() !== endedSession ||
this.player.getState() !== "idle" this.player.getState() !== "idle"
) return true; ) {
if (this.streamRecovery === state) this.streamRecovery = null;
return true;
}
if (!result?.url || !this.connected) return false; if (!result?.url || !this.connected) return false;
song.url = result.url; song.url = result.url;
this.player.play(result.url, position, duration); 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"); this.emit("stateChange");
return true; return true;
} }
@@ -1438,6 +1454,10 @@ export class BotInstance extends EventEmitter {
} }
private cmdPause(): string { private cmdPause(): string {
const recovery = this.streamRecovery;
if (recovery && recovery.song === this.queue.current() && this.player.getState() === "idle") {
recovery.pauseRequested = true;
}
this.player.pause(); this.player.pause();
if (this.queue.current()?.platform === "spotify") { if (this.queue.current()?.platform === "spotify") {
this.spotifyController.pause().catch((err) => this.spotifyController.pause().catch((err) =>
@@ -1450,7 +1470,15 @@ export class BotInstance extends EventEmitter {
} }
private cmdResume(): string { private cmdResume(): string {
const recovery = this.streamRecovery;
const retryInterrupted = recovery && recovery.song === this.queue.current() &&
recovery.session === this.player.getPlaybackSessionId() && !recovery.inFlight &&
recovery.pauseRequested && this.player.getState() === "idle";
if (recovery) recovery.pauseRequested = false;
this.player.resume(); this.player.resume();
// A lookup that failed while paused has no stream to resume. Re-enter
// the bounded end/recovery handler instead of reporting success forever idle.
if (retryInterrupted) this.player.emit("trackEnd");
if (this.queue.current()?.platform === "spotify") { if (this.queue.current()?.platform === "spotify") {
this.spotifyController.resume().catch((err) => this.spotifyController.resume().catch((err) =>
this.logger.warn({ err }, "Spotify resume failed")); this.logger.warn({ err }, "Spotify resume failed"));