feat(audio): add external-PCM mode (playPcmStream) for Spotify sidecar

Adds AudioPlayer.playPcmStream(readable, {onExternalEnd}) that feeds a
long-lived external 48kHz/s16le/stereo Readable into the existing pcmBuffer +
20ms frame loop + Opus encoder without spawning a per-URL ffmpeg. Reuses the
same high/low-water backpressure (pausing/resuming the Readable), suppresses the
underrun trackEnd drain/stall branches while external (emitting a silence frame
to keep the 20ms timeline), tears down externalMode in stop() by DETACHING the
shared readable (remove our data/end/error listeners + pause, never destroy the
sidecar stream) and clearing onExternalEnd, and makes seek() a local no-op in
external mode. The url play() path and all exported pure functions are unchanged.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
saopig1andClaude Opus 4.8 committed 2026-07-02 21:14:15 +08:00
1 parent 641da086e5
commit 6debe23034
2 files changed
+355 -8

No files matched your search

+192 -1
View File
@@ -2,7 +2,9 @@ import { describe, it, expect } from "vitest";
import { mkdtempSync, writeFileSync, existsSync } from "node:fs"; import { mkdtempSync, writeFileSync, existsSync } from "node:fs";
import { tmpdir } from "node:os"; import { tmpdir } from "node:os";
import { join } from "node:path"; import { join } from "node:path";
import { buildFfmpegArgs, shouldUsePowerShellDownload, cleanupTempDir, shouldEndOnStall, volumeToFactor } from "./player.js"; import { Readable } from "node:stream";
import { buildFfmpegArgs, shouldUsePowerShellDownload, cleanupTempDir, shouldEndOnStall, volumeToFactor, AudioPlayer } from "./player.js";
import type { Logger } from "../logger.js";
function getHeadersArg(args: string[]): string { function getHeadersArg(args: string[]): string {
const idx = args.indexOf("-headers"); const idx = args.indexOf("-headers");
@@ -200,3 +202,192 @@ describe("shouldEndOnStall (#89 mid-track stall watchdog)", () => {
expect(shouldEndOnStall(10, false, MAX_EMPTY, MAX_STALL)).toBe(false); expect(shouldEndOnStall(10, false, MAX_EMPTY, MAX_STALL)).toBe(false);
}); });
}); });
// Minimal stub: AudioPlayer only calls debug/info/warn/error; child() returns self.
const silentLogger = {
debug() {},
info() {},
warn() {},
error() {},
fatal() {},
trace() {},
child() {
return silentLogger;
},
} as unknown as Logger;
// A readable we fully control: no underlying source; we push PCM manually and
// keep it open (never push(null)) to model the long-lived go-librespot sidecar.
function openPcmReadable(): Readable {
return new Readable({ read() {} });
}
const wait = (ms: number): Promise<void> => new Promise<void>((r) => setTimeout(r, ms));
const FRAME_BYTES = 3840; // PCM_FRAME_BYTES: 960 samples * 2ch * 2 bytes @48k s16le
describe("AudioPlayer external-PCM mode (playPcmStream)", () => {
it("emits Opus 'frame' events from the external PCM stream without spawning ffmpeg", async () => {
const player = new AudioPlayer(silentLogger);
const frames: Buffer[] = [];
player.on("frame", (f) => frames.push(f));
const stream = openPcmReadable();
player.playPcmStream(stream, {});
stream.push(Buffer.alloc(FRAME_BYTES * 10)); // ~10 frames of PCM
await wait(150); // ~7 frame ticks at 20ms
expect(player.getState()).toBe("playing");
expect(frames.length).toBeGreaterThan(0);
expect(Buffer.isBuffer(frames[0])).toBe(true);
player.stop();
});
it("does NOT emit 'trackEnd' on underrun while external (stream stays open)", async () => {
const player = new AudioPlayer(silentLogger);
let ended = 0;
const frames: Buffer[] = [];
player.on("trackEnd", () => ended++);
player.on("frame", (f) => frames.push(f));
const stream = openPcmReadable();
player.playPcmStream(stream, {});
stream.push(Buffer.alloc(FRAME_BYTES * 2)); // only 2 frames, then underrun
await wait(200); // long after those 2 frames have drained
// In the url path, ffmpeg===null + empty buffer would fire trackEnd; here it must not.
expect(ended).toBe(0);
// Silence frames keep the 20ms timeline alive -> more than the 2 fed frames emitted.
expect(frames.length).toBeGreaterThan(2);
expect(player.getState()).toBe("playing");
player.stop();
});
// CORRECTION C2 (c): stop() DETACHES the shared readable — it must NOT be destroyed
// (destroying the sidecar's long-lived ffmpeg stdout would kill it for every future
// track). The sessionId bump + listener removal fence stale PCM out of pcmBuffer.
it("stop() detaches external mode without destroying the readable, and fences via sessionId", async () => {
const player = new AudioPlayer(silentLogger);
const frames: Buffer[] = [];
player.on("frame", (f) => frames.push(f));
const stream = openPcmReadable();
player.playPcmStream(stream, {});
expect(stream.listenerCount("data")).toBe(1);
stream.push(Buffer.alloc(FRAME_BYTES * 5));
await wait(80);
player.stop();
expect(player.getState()).toBe("idle");
// C2: the shared sidecar stream must NOT be destroyed by teardown.
expect(stream.destroyed).toBe(false);
// Player's listeners are removed on detach (data/end/error).
expect(stream.listenerCount("data")).toBe(0);
expect(stream.listenerCount("end")).toBe(0);
expect(stream.listenerCount("error")).toBe(0);
const countAtStop = frames.length;
// sessionId fence + detached listeners: PCM pushed after stop must not
// resurrect the timeline or re-feed pcmBuffer.
stream.push(Buffer.alloc(FRAME_BYTES * 5));
await wait(80);
expect(frames.length).toBe(countAtStop);
});
// CORRECTION C2 (a): a gapless track change is driven by the sidecar pushing LATER
// PCM over the SAME already-attached stream. The player must NOT detach/re-attach
// (no second playPcmStream) — one persistent data listener serves every track.
it("(C2-a) feeds a later chunk over the SAME single attachment — gapless track change, no re-attach", async () => {
const player = new AudioPlayer(silentLogger);
const frames: Buffer[] = [];
player.on("frame", (f) => frames.push(f));
const stream = openPcmReadable();
player.playPcmStream(stream, {});
expect(stream.listenerCount("data")).toBe(1); // attached exactly once
stream.push(Buffer.alloc(FRAME_BYTES * 4)); // "track 1" PCM
await wait(120);
const afterFirst = frames.length;
expect(afterFirst).toBeGreaterThan(0);
stream.push(Buffer.alloc(FRAME_BYTES * 4)); // sidecar seamlessly rolls into "track 2"
await wait(120);
expect(frames.length).toBeGreaterThan(afterFirst);
// Still exactly ONE listener — no detach/re-attach across the handoff.
expect(stream.listenerCount("data")).toBe(1);
expect(player.getState()).toBe("playing");
player.stop();
});
// CORRECTION C2 (b): a second playPcmStream detaches the first (NOT destroyed, and it
// stops feeding pcmBuffer) and attaches the second.
it("(C2-b) a second playPcmStream detaches the first (not destroyed, stops feeding) and attaches the second", async () => {
const player = new AudioPlayer(silentLogger);
const frames: Buffer[] = [];
player.on("frame", (f) => frames.push(f));
const first = openPcmReadable();
player.playPcmStream(first, {});
first.push(Buffer.alloc(FRAME_BYTES * 4));
await wait(120);
expect(frames.length).toBeGreaterThan(0);
expect(first.listenerCount("data")).toBe(1);
const second = openPcmReadable();
player.playPcmStream(second, {}); // fences + detaches `first`, attaches `second`
// C2: `first` is DETACHED, not destroyed.
expect(first.destroyed).toBe(false);
// `first` no longer feeds pcmBuffer — its data listener was removed.
expect(first.listenerCount("data")).toBe(0);
// `second` is now the attached source.
expect(second.listenerCount("data")).toBe(1);
expect(player.getState()).toBe("playing");
player.stop();
});
it("fires onExternalEnd when the readable ends (drives controller-based advance)", async () => {
const player = new AudioPlayer(silentLogger);
let endedCb = 0;
const stream = openPcmReadable();
player.playPcmStream(stream, { onExternalEnd: () => endedCb++ });
stream.push(Buffer.alloc(FRAME_BYTES));
await wait(40);
stream.push(null); // end-of-stream
await wait(40);
expect(endedCb).toBe(1);
player.stop();
});
it("seek() is a local no-op in external mode (never respawns ffmpeg on a spotify sentinel)", async () => {
const player = new AudioPlayer(silentLogger);
const stream = openPcmReadable();
player.playPcmStream(stream, {});
stream.push(Buffer.alloc(FRAME_BYTES * 3));
await wait(40);
expect(() => player.seek(30)).not.toThrow();
// Still external, still playing — no url-ffmpeg respawn, state unchanged.
expect(player.getState()).toBe("playing");
player.stop();
});
it("pause()/resume() still gate local emission in external mode (unchanged semantics)", async () => {
const player = new AudioPlayer(silentLogger);
const stream = openPcmReadable();
player.playPcmStream(stream, {});
stream.push(Buffer.alloc(FRAME_BYTES * 3));
await wait(40);
player.pause();
expect(player.getState()).toBe("paused");
player.resume();
expect(player.getState()).toBe("playing");
player.stop();
});
});
+163 -7
View File
@@ -5,6 +5,7 @@ import { accessSync, chmodSync, constants, mkdtempSync, rmSync } from "node:fs";
import { tmpdir } from "node:os"; import { tmpdir } from "node:os";
import { join } from "node:path"; import { join } from "node:path";
import { createOpusEncoder, PCM_FRAME_BYTES, type Encoder } from "./encoder.js"; import { createOpusEncoder, PCM_FRAME_BYTES, type Encoder } from "./encoder.js";
import type { Readable } from "node:stream";
import type { Logger } from "../logger.js"; import type { Logger } from "../logger.js";
const require = createRequire(import.meta.url); const require = createRequire(import.meta.url);
@@ -187,6 +188,22 @@ export class AudioPlayer extends EventEmitter {
private static readonly MAX_STALL_ATTEMPTS = 3000; private static readonly MAX_STALL_ATTEMPTS = 3000;
private currentSongDuration = 0; // 当前歌曲总时长(秒) private currentSongDuration = 0; // 当前歌曲总时长(秒)
// --- External PCM mode (Stage 2: go-librespot Spotify sidecar) ---
// When true, PCM arrives from a long-lived external Readable instead of a
// per-URL ffmpeg: this.ffmpeg stays null, and the underrun-driven trackEnd
// branches are suppressed (advance is driven by the controller, not EOF).
//
// CORRECTION C2: externalStream is the backend's LONG-LIVED, SHARED ffmpeg
// stdout (one stream reused across every track). Teardown must DETACH (remove
// the listeners we added + pause), never destroy it. We keep references to the
// exact handler functions so detach can removeListener precisely.
private externalMode = false;
private externalStream: Readable | null = null;
private onExternalEnd: (() => void) | null = null;
private externalDataHandler: ((chunk: Buffer) => void) | null = null;
private externalEndHandler: (() => void) | null = null;
private externalErrorHandler: ((err: Error) => void) | null = null;
constructor(logger: Logger) { constructor(logger: Logger) {
super(); super();
this.encoder = createOpusEncoder(); this.encoder = createOpusEncoder();
@@ -390,6 +407,108 @@ export class AudioPlayer extends EventEmitter {
this.startFrameLoop(); this.startFrameLoop();
} }
/**
* External-PCM mode (Stage 2 go-librespot Spotify sidecar).
*
* Feeds an already-normalized 48kHz/s16le/stereo PCM Readable (the
* go-librespot FIFO -> ffmpeg output) straight into the existing pcmBuffer +
* 20ms frame loop + Opus encoder + "frame" emission, WITHOUT spawning a
* per-URL ffmpeg. The url play() path is left completely untouched.
*
* Track advance is NOT driven by buffer underrun here (the sidecar stream is
* continuous and never EOFs per song); the caller drives advance via the
* SpotifyController "trackEnded" WebSocket event. onExternalEnd fires only if
* the underlying readable itself ends or errors.
*
* CORRECTION C2: the readable is the backend's long-lived, SHARED ffmpeg
* stdout reused across every track — a gapless track change is just LATER PCM
* on this SAME already-attached stream (no re-attach). Teardown DETACHES
* (removes our listeners + pauses); it never destroys the shared stream.
*/
playPcmStream(readable: Readable, opts: { onExternalEnd?: () => void } = {}): void {
// 1. Fence current playback: stop() bumps sessionId, clears pcmBuffer, kills
// any ffmpeg, and DETACHES (never destroys) any prior external stream.
this.stop();
const currentSessionId = this.sessionId;
this.externalMode = true;
this.externalStream = readable;
this.onExternalEnd = opts.onExternalEnd ?? null;
// Leave this.ffmpeg = null; clear currentUrl so seek() cannot respawn ffmpeg.
this.currentUrl = "";
this.seekOffset = 0;
this.framesPlayed = 0;
this.healthyFrames = 0;
this.ffmpegPaused = false;
this.spawnFailed = false;
this.emptyFrameAttempts = 0;
this.currentSongDuration = 0;
// Same ingestion + high-water backpressure as the ffmpeg.stdout handler,
// but pausing the Readable instead of ffmpeg.stdout. sessionId-guarded so
// stale sidecar PCM can't leak into a new track after stop()/skip. Handler
// refs are stored so detach can remove exactly these listeners (C2).
const onData = (chunk: Buffer): void => {
if (this.sessionId !== currentSessionId) return;
this.pcmBuffer = Buffer.concat([this.pcmBuffer, chunk]);
if (
this.pcmBuffer.length > AudioPlayer.BUFFER_HIGH_WATER &&
!this.ffmpegPaused &&
this.externalStream === readable
) {
readable.pause();
this.ffmpegPaused = true;
}
};
const onEnd = (): void => {
if (this.sessionId !== currentSessionId) return;
this.onExternalEnd?.();
};
const onError = (err: Error): void => {
if (this.sessionId !== currentSessionId) return;
this.logger.warn({ err }, "External PCM stream error");
this.onExternalEnd?.();
};
this.externalDataHandler = onData;
this.externalEndHandler = onEnd;
this.externalErrorHandler = onError;
readable.on("data", onData);
readable.on("end", onEnd);
readable.on("error", onError);
this.state = "playing";
this.startFrameLoop();
}
/**
* CORRECTION C2: DETACH, never destroy. The external readable is the backend's
* long-lived, SHARED ffmpeg stdout reused across every track; destroying it
* would kill the sidecar pipe for all future tracks. Remove only the listeners
* WE added and pause the flow so stale PCM stops landing in pcmBuffer, then
* clear the external-mode state.
*/
private detachExternalStream(): void {
const stream = this.externalStream;
if (stream) {
if (this.externalDataHandler) stream.off("data", this.externalDataHandler);
if (this.externalEndHandler) stream.off("end", this.externalEndHandler);
if (this.externalErrorHandler) stream.off("error", this.externalErrorHandler);
try {
stream.pause();
} catch {
/* best-effort: never destroy the shared sidecar stream */
}
}
this.externalDataHandler = null;
this.externalEndHandler = null;
this.externalErrorHandler = null;
this.externalStream = null;
this.externalMode = false;
this.onExternalEnd = null;
}
stop(): void { stop(): void {
// 3. 递增 ID 是最有效的逻辑“隔离墙” // 3. 递增 ID 是最有效的逻辑“隔离墙”
this.sessionId++; this.sessionId++;
@@ -419,6 +538,11 @@ export class AudioPlayer extends EventEmitter {
this.currentTempDir = null; this.currentTempDir = null;
} }
// CORRECTION C2: tear down external mode by DETACHING (remove our listeners +
// pause) — never destroy the shared, long-lived sidecar stream. The
// sessionId++ above already fences the external data/end/error handlers.
this.detachExternalStream();
this.ffmpegPaused = false; this.ffmpegPaused = false;
this.spawnFailed = false; this.spawnFailed = false;
this.state = "idle"; this.state = "idle";
@@ -480,7 +604,10 @@ export class AudioPlayer extends EventEmitter {
? (this.currentSongDuration - elapsed) <= 5 // 距离结尾不足5秒 ? (this.currentSongDuration - elapsed) <= 5 // 距离结尾不足5秒
: true; // 未知时长时保守处理 : true; // 未知时长时保守处理
if (this.ffmpeg !== null && this.pcmBuffer.length < PCM_FRAME_BYTES) { // External mode: the sidecar PCM stream is continuous and never EOFs per
// song; a transient underrun must NOT end the track (advance is driven by
// the controller). Skip BOTH drain/stall branches while externalMode.
if (!this.externalMode && this.ffmpeg !== null && this.pcmBuffer.length < PCM_FRAME_BYTES) {
this.emptyFrameAttempts++; this.emptyFrameAttempts++;
// End the track when FFmpeg has gone silent: quickly if we're near the // End the track when FFmpeg has gone silent: quickly if we're near the
@@ -526,7 +653,7 @@ export class AudioPlayer extends EventEmitter {
this.emptyFrameAttempts = 0; this.emptyFrameAttempts = 0;
} }
if (!this.ffmpeg && this.pcmBuffer.length < PCM_FRAME_BYTES) { if (!this.externalMode && !this.ffmpeg && this.pcmBuffer.length < PCM_FRAME_BYTES) {
this.frameLoopRunning = false; this.frameLoopRunning = false;
if (this.state !== "idle") { if (this.state !== "idle") {
this.state = "idle"; this.state = "idle";
@@ -542,13 +669,24 @@ export class AudioPlayer extends EventEmitter {
} }
private sendNextFrame(): void { private sendNextFrame(): void {
if (this.pcmBuffer.length < PCM_FRAME_BYTES) return; if (this.pcmBuffer.length < PCM_FRAME_BYTES) {
// External mode: the sidecar PCM stream is long-lived and must NOT end on
// a transient underrun. Emit an encoded silence frame so the 20ms voice
// timeline stays continuous instead of returning (which would desync TS).
if (this.externalMode) this.emitSilenceFrame();
return;
}
const pcmFrame = this.pcmBuffer.subarray(0, PCM_FRAME_BYTES); const pcmFrame = this.pcmBuffer.subarray(0, PCM_FRAME_BYTES);
this.pcmBuffer = this.pcmBuffer.subarray(PCM_FRAME_BYTES); this.pcmBuffer = this.pcmBuffer.subarray(PCM_FRAME_BYTES);
if (this.ffmpegPaused && this.pcmBuffer.length < AudioPlayer.BUFFER_LOW_WATER && this.ffmpeg?.stdout) { if (this.ffmpegPaused && this.pcmBuffer.length < AudioPlayer.BUFFER_LOW_WATER) {
this.ffmpeg.stdout.resume(); if (this.externalMode && this.externalStream) {
this.ffmpegPaused = false; this.externalStream.resume();
this.ffmpegPaused = false;
} else if (this.ffmpeg?.stdout) {
this.ffmpeg.stdout.resume();
this.ffmpegPaused = false;
}
} }
try { try {
@@ -566,6 +704,16 @@ export class AudioPlayer extends EventEmitter {
} }
} }
private emitSilenceFrame(): void {
try {
const opusFrame = this.encoder.encode(Buffer.alloc(PCM_FRAME_BYTES));
this.emit("frame", opusFrame);
this.framesPlayed++;
} catch (err) {
this.emit("error", err as Error);
}
}
private applyVolume(pcm: Buffer): Buffer { private applyVolume(pcm: Buffer): Buffer {
const factor = volumeToFactor(this.volume); const factor = volumeToFactor(this.volume);
// factor === 1 only at volume 100; skip the per-sample loop at full loudness. // factor === 1 only at volume 100; skip the per-sample loop at full loudness.
@@ -578,8 +726,16 @@ export class AudioPlayer extends EventEmitter {
return out; return out;
} }
// NOTE: in external (Spotify sidecar) mode getElapsed() is frame-count based
// (framesPlayed includes silence frames emitted on underrun) and therefore
// only APPROXIMATE — the authoritative position is the controller's live
// status.track.position. This approximation is acceptable for Spotify.
getElapsed(): number { return this.seekOffset + (this.framesPlayed * FRAME_DURATION_MS) / 1000; } getElapsed(): number { return this.seekOffset + (this.framesPlayed * FRAME_DURATION_MS) / 1000; }
seek(seconds: number): void { seek(seconds: number): void {
// External (Spotify sidecar) mode: local seek is a no-op. Respawning ffmpeg
// on the spotify: sentinel would collide with the continuous PCM source;
// transport is delegated to the SpotifyController by the caller (Task 7).
if (this.externalMode) return;
if (this.currentUrl && Number.isFinite(seconds) && seconds >= 0) { if (this.currentUrl && Number.isFinite(seconds) && seconds >= 0) {
this.play(this.currentUrl, seconds, this.currentSongDuration); this.play(this.currentUrl, seconds, this.currentSongDuration);
} }