fix(spotify): don't skip paused Rust track; handle ffmpeg stdin EPIPE; guard sub-window end-detection [corner-case]

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
saopig1andClaude Opus 4.8 committed 2026-07-03 11:18:47 +08:00
1 parent 11c0948330
commit 6dc99d88e2
2 files changed
+182 -5

No files matched your search

+105 -3
View File
@@ -205,7 +205,41 @@ describe("RustLibrespotBackend track-end poll loop", () => {
expect(ended).toHaveBeenCalledWith({ uri: "spotify:track:A", reason: "ended" });
});
it("emits trackEnded once when playback stops (!isPlaying) after having played", async () => {
// C1(pause-skip): a USER pause reports is_playing:false with the SAME uri on
// the Rust backend. That MUST NOT be read as a track end (it would skip the
// paused track and break pause + occupancy auto-pause). Formerly the
// "!isPlaying after having played" test asserted the opposite — that encoded
// the bug; it is now split into this pause-no-skip test plus the two-poll
// external-stop test below.
it("does NOT emit trackEnded when the user PAUSES (self-initiated pause is not a track end)", async () => {
const h = makeHarness();
const ended = vi.fn();
h.backend.on("trackEnded", ended);
await h.backend.playTrack("spotify:track:A");
// Observe our track actually playing first.
h.connect.getPlaybackState.mockResolvedValueOnce({
isPlaying: true, progressMs: 5000, trackUri: "spotify:track:A", durationMs: 200000,
});
await (h.backend as any).pollState();
// User pauses: the Connect device stays loaded but reports is_playing:false
// with the SAME uri across every subsequent poll while paused.
await h.backend.pause();
h.connect.getPlaybackState.mockResolvedValue({
isPlaying: false, progressMs: 5000, trackUri: "spotify:track:A", durationMs: 200000,
});
await (h.backend as any).pollState();
await (h.backend as any).pollState(); // stays paused across multiple polls
expect(ended).not.toHaveBeenCalled();
// Resuming keeps the same track playing — still no spurious end.
await h.backend.resume();
h.connect.getPlaybackState.mockResolvedValue({
isPlaying: true, progressMs: 6000, trackUri: "spotify:track:A", durationMs: 200000,
});
await (h.backend as any).pollState();
expect(ended).not.toHaveBeenCalled();
});
it("emits trackEnded once on an EXTERNAL stop only after TWO consecutive !isPlaying polls (a transient mid-track !isPlaying is not a skip)", async () => {
const h = makeHarness();
const ended = vi.fn();
h.backend.on("trackEnded", ended);
@@ -214,13 +248,45 @@ describe("RustLibrespotBackend track-end poll loop", () => {
.mockResolvedValueOnce({ isPlaying: true, progressMs: 5000, trackUri: "spotify:track:A", durationMs: 200000 })
.mockResolvedValueOnce({ isPlaying: false, progressMs: 5000, trackUri: "spotify:track:A", durationMs: 200000 })
.mockResolvedValue({ isPlaying: false, progressMs: 5000, trackUri: "spotify:track:A", durationMs: 200000 });
await (h.backend as any).pollState();
await (h.backend as any).pollState();
await (h.backend as any).pollState(); // observed playing
await (h.backend as any).pollState(); // FIRST !isPlaying -> unconfirmed (could be transient buffering)
expect(ended).not.toHaveBeenCalled();
await (h.backend as any).pollState(); // SECOND consecutive !isPlaying -> confirmed external stop
await (h.backend as any).pollState(); // idempotent: no second emit for same track
expect(ended).toHaveBeenCalledTimes(1);
expect(ended).toHaveBeenCalledWith({ uri: "spotify:track:A", reason: "ended" });
});
it("a transient single !isPlaying poll followed by playing again does NOT emit trackEnded (buffering hiccup)", async () => {
const h = makeHarness();
const ended = vi.fn();
h.backend.on("trackEnded", ended);
await h.backend.playTrack("spotify:track:A");
h.connect.getPlaybackState
.mockResolvedValueOnce({ isPlaying: true, progressMs: 5000, trackUri: "spotify:track:A", durationMs: 200000 })
.mockResolvedValueOnce({ isPlaying: false, progressMs: 5000, trackUri: "spotify:track:A", durationMs: 200000 })
.mockResolvedValue({ isPlaying: true, progressMs: 6000, trackUri: "spotify:track:A", durationMs: 200000 });
await (h.backend as any).pollState(); // playing
await (h.backend as any).pollState(); // momentary !isPlaying (buffering)
await (h.backend as any).pollState(); // playing again -> stop confirmation reset
expect(ended).not.toHaveBeenCalled();
});
// m(sub-window): a track SHORTER than the end-of-track window must not be
// declared finished on its first observed-playing poll (durationMs - window
// is negative, so the old near-end check fired unconditionally).
it("does NOT false-finish a sub-window (< END_OF_TRACK_WINDOW_MS) duration on the first playing poll", async () => {
const h = makeHarness();
const ended = vi.fn();
h.backend.on("trackEnded", ended);
await h.backend.playTrack("spotify:track:short");
h.connect.getPlaybackState.mockResolvedValue({
isPlaying: true, progressMs: 100, trackUri: "spotify:track:short", durationMs: 1200,
});
await (h.backend as any).pollState();
expect(ended).not.toHaveBeenCalled();
});
it("emits trackEnded when the track uri transitions to null after playing", async () => {
const h = makeHarness();
const ended = vi.fn();
@@ -437,4 +503,40 @@ describe("RustLibrespotBackend child-process error handling", () => {
expect(onErr).toHaveBeenCalledWith(err);
h.backend.stop();
});
// I(pipe): ffmpeg dying mid-track while librespot keeps producing PCM raises
// EPIPE on ffmpeg.stdin. With no stdin 'error' listener Node escalates it to
// process 'uncaughtException'. The backend must handle it in-band.
it("swallows an EPIPE 'error' on ffmpeg.stdin (ffmpeg died mid-track) without an unhandled throw", async () => {
const h = makeHarness();
await h.backend.start();
expect(h.backend.listenerCount("error")).toBe(0);
const epipe = Object.assign(new Error("write EPIPE"), { code: "EPIPE" });
expect(() => h.ffmpegChild.stdin.emit("error", epipe)).not.toThrow();
h.backend.stop(); // idempotent second teardown must not throw
});
it("routes an ffmpeg.stdin EPIPE to the backend 'error' listener and tears down cleanly", async () => {
const h = makeHarness();
await h.backend.start();
const onErr = vi.fn();
h.backend.on("error", onErr);
const epipe = Object.assign(new Error("write EPIPE"), { code: "EPIPE" });
expect(() => h.ffmpegChild.stdin.emit("error", epipe)).not.toThrow();
expect(onErr).toHaveBeenCalledWith(epipe);
// Broken pipe -> clean teardown: children killed, not ready.
expect(h.librespotChild.kill).toHaveBeenCalled();
expect(h.ffmpegChild.kill).toHaveBeenCalled();
expect(h.backend.isReady()).toBe(false);
h.backend.stop(); // second teardown must not throw (no double-teardown crash)
expect(onErr).toHaveBeenCalledTimes(1); // single emit despite both pipe ends
});
it("does not throw when librespot proc.stdout emits an EPIPE on the broken pipe", async () => {
const h = makeHarness();
await h.backend.start();
const epipe = Object.assign(new Error("read/write EPIPE"), { code: "EPIPE" });
expect(() => h.librespotChild.stdout.emit("error", epipe)).not.toThrow();
h.backend.stop();
});
});
+77 -2
View File
@@ -91,6 +91,21 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
// side effects at all, so a foreign track already near its end can't emit a
// spurious trackEnded/metadata and wrongly advance the queue.
private armed = false;
// C1(pause-skip): our own pause state. A self-initiated pause makes the
// Connect device report is_playing:false with the SAME uri; without tracking
// it we'd misread that as a track end and SKIP the paused track (breaking the
// pause command and occupancy auto-pause-when-alone on the Rust backend). Set
// by pause(), cleared by resume() and playTrack().
private paused = false;
// Robustness: a genuine external stop must be confirmed across TWO consecutive
// non-paused is_playing:false polls before we emit trackEnded, so a momentary
// mid-track is_playing:false (buffering) can't false-skip. Set on the first
// such poll; reset whenever the track is playing again or the track changes.
private stopSeen = false;
// I(pipe): one teardown per broken librespot->ffmpeg pipe. Both stream ends
// (ffmpeg.stdin write side, proc.stdout read side) can report the same EPIPE;
// this guard prevents a double teardown / double error-emit.
private pipeBroke = false;
constructor(o: RustLibrespotBackendOptions) {
super();
@@ -121,6 +136,7 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
// Everything past here spawns children / opens the state poll. On any
// failure (e.g. the device never appears), tear it all down via stop().
try {
this.pipeBroke = false; // fresh pipe for this start()
// 1. Spawn ffmpeg FIRST (the reader) so its stdin pipe is ready before
// librespot starts pushing raw 44.1k s16le PCM into it.
this.ffmpeg = spawn(
@@ -135,6 +151,12 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
this.log.debug({ ffmpeg: b.toString().trim() }, "ffmpeg"),
);
this.ffmpeg.on("error", (err) => this.emitError(err));
// I(pipe): if ffmpeg dies mid-track while librespot keeps producing PCM,
// the next write into its stdin hits the now-closed pipe -> EPIPE 'error'
// on stdin. With no listener Node escalates that to process
// 'uncaughtException' (the global handler logs but leaves undefined
// state). Handle it in-band: tear down cleanly and surface via emitError.
this.ffmpeg.stdin?.on("error", (err) => this.onPipeError(err));
// 2. Spawn librespot: --backend pipe with NO --device => raw s16le/44100/2
// on stdout, NO --passthrough (that would emit raw Ogg). --access-token
@@ -160,6 +182,10 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
// stdout carries PCM — pipe it, never attach a data listener that consumes it.
if (this.proc.stdout && this.ffmpeg.stdin) {
this.proc.stdout.pipe(this.ffmpeg.stdin);
// Guard the read side too: a broken pipe can surface as an 'error' on
// proc.stdout when ffmpeg's stdin closes underneath it. Same handler,
// guarded against a double teardown.
this.proc.stdout.on("error", (err) => this.onPipeError(err));
}
this.proc.stderr?.on("data", (b: Buffer) =>
this.log.info({ librespot: b.toString().trim() }, "librespot"),
@@ -196,6 +222,25 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
}
}
/**
* I(pipe): the librespot->ffmpeg PCM pipe broke — typically ffmpeg died
* mid-track and librespot's next write hit the closed pipe (EPIPE), or the
* stream was destroyed (ERR_STREAM_DESTROYED). Both are expected teardown
* signals, not programming errors. Log, tear down cleanly (stop() is
* idempotent), then surface via emitError so a listening controller reacts.
* Guarded so both stream ends reporting the same break don't double-teardown.
*/
private onPipeError(err: unknown): void {
if (this.pipeBroke) return;
this.pipeBroke = true;
this.log.warn(
{ err },
"rust-librespot: librespot->ffmpeg pipe broke (ffmpeg died?) — tearing down",
);
this.stop();
this.emitError(err);
}
private async waitForDevice(): Promise<void> {
const sleep = this.deps.sleep ?? defaultSleep;
const interval = this.deps.readyPollIntervalMs ?? DEFAULT_READY_POLL_MS;
@@ -265,6 +310,7 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
this.currentUri = state.trackUri;
this.hasPlayed = false;
this.endedForCurrent = false;
this.stopSeen = false; // new track -> drop any pending stop confirmation
const np: SpotifyNowPlaying = {
uri: state.trackUri,
name: "",
@@ -278,6 +324,8 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
if (state.isPlaying) {
this.hasPlayed = true;
// Playing again -> any earlier is_playing:false was transient, not a stop.
this.stopSeen = false;
// Real playback observed -> the I4 degrade-to-skip watchdog is moot.
this.clearPlaybackWatchdog();
}
@@ -287,11 +335,29 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
// the bot's own track has actually been observed playing. Without this, the
// first poll (before playTrack) could observe the account's stale/paused
// track sitting near its end and spuriously emit "trackEnded".
// m(sub-window): only apply the near-end window when the track is LONGER
// than the window. For durationMs in [1, END_OF_TRACK_WINDOW_MS],
// `durationMs - window` is negative, so the old `> 0` guard made this
// unconditionally true and finished the track on its first poll. A
// sub-window track instead relies on normal stop/next-track detection.
const finishedByProgress =
this.hasPlayed &&
state.durationMs > 0 &&
state.durationMs > END_OF_TRACK_WINDOW_MS &&
state.progressMs >= state.durationMs - END_OF_TRACK_WINDOW_MS;
const finishedByStop = this.hasPlayed && !state.isPlaying;
// C1(pause-skip): a self-initiated pause (this.paused) reports is_playing:false
// with the SAME uri — that is NOT a track end, so never finish while paused.
// And even for a genuine external stop, require it to persist across two
// consecutive polls (stopSeen) so a momentary mid-track is_playing:false
// (buffering) doesn't false-skip. Confirmed stops still emit within ~one
// extra poll interval.
let finishedByStop = false;
if (this.hasPlayed && !state.isPlaying && !this.paused) {
if (this.stopSeen) {
finishedByStop = true;
} else {
this.stopSeen = true; // first non-paused stop poll — await confirmation
}
}
const finishedByNull = this.hasPlayed && state.trackUri === null;
if (finishedByProgress || finishedByStop || finishedByNull) {
@@ -323,6 +389,10 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
this.hasPlayed = false;
this.endedForCurrent = false;
this.armed = true;
// A newly-started track is not paused, and carries no pending stop
// confirmation from the previous track. (C1 pause-skip.)
this.paused = false;
this.stopSeen = false;
// transfer(false) activates our device WITHOUT starting audio; play() then
// actually starts the uri. The two-step is required — transfer alone won't
// begin playback.
@@ -384,10 +454,15 @@ export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBa
async pause(): Promise<void> {
await this.connect.pause();
// C1(pause-skip): mark our own pause so the next poll's is_playing:false
// (same uri) is not misread as a track end and skipped.
this.paused = true;
}
async resume(): Promise<void> {
await this.connect.resume();
// Resumed -> normal end-detection applies again.
this.paused = false;
}
async seek(ms: number): Promise<void> {