diff --git a/src/music/spotify/rust-librespot.test.ts b/src/music/spotify/rust-librespot.test.ts new file mode 100644 index 0000000..c155030 --- /dev/null +++ b/src/music/spotify/rust-librespot.test.ts @@ -0,0 +1,283 @@ +import { describe, it, expect, vi } from "vitest"; +import { EventEmitter } from "node:events"; +import { PassThrough } from "node:stream"; +import pino from "pino"; +import { RustLibrespotBackend } from "./rust-librespot.js"; + +const log = pino({ level: "silent" }); + +/** ChildProcess stand-in with real Readable/Writable pipes so stdout->stdin piping works. */ +function makeFakeChild() { + const child: any = new EventEmitter(); + child.stdout = new PassThrough(); + child.stderr = new PassThrough(); + child.stdin = new PassThrough(); + child.kill = vi.fn(); + return child; +} + +function makeConnect() { + return { + getDevices: vi.fn(async () => [{ id: "dev1", name: "Test Bot", is_active: false }]), + findDeviceByName: vi.fn(async () => "dev1"), + transfer: vi.fn(async () => {}), + play: vi.fn(async () => {}), + pause: vi.fn(async () => {}), + resume: vi.fn(async () => {}), + seek: vi.fn(async () => {}), + getPlaybackState: vi.fn(async () => null as any), + }; +} + +function makeOAuth() { + return { + getAccessToken: vi.fn(async () => "tok-123" as string | null), + isAuthorized: () => true, + }; +} + +function makeHarness(over: { connect?: any; oauth?: any } = {}) { + const calls: string[] = []; + const librespotChild = makeFakeChild(); + const ffmpegChild = makeFakeChild(); + + const spawn = vi.fn((cmd: string, ..._rest: any[]) => { + const isLibrespot = cmd.includes("librespot"); + calls.push(`spawn:${isLibrespot ? "librespot" : cmd}`); + return isLibrespot ? librespotChild : ffmpegChild; + }); + const mkdirSync = vi.fn(); + const connect = over.connect ?? makeConnect(); + const oauth = over.oauth ?? makeOAuth(); + + const backend = new RustLibrespotBackend({ + deviceName: "Test Bot", + bitrate: 320, + cacheDir: "/tmp/cache", + oauth: oauth as any, + connect: connect as any, + logger: log, + deps: { + spawn: spawn as any, + mkdirSync: mkdirSync as any, + findBinary: () => "/bin/librespot", + // C1: pin ffmpeg so arg-array assertions stay stable while prod uses getFfmpegCommand(). + ffmpegCommand: "ffmpeg", + sleep: async () => {}, + readyPollIntervalMs: 1, + readyTimeoutMs: 100, + // huge so the background setInterval never fires; tests drive pollState() directly. + statePollIntervalMs: 10_000_000, + }, + }); + + return { backend, calls, spawn, mkdirSync, connect, oauth, librespotChild, ffmpegChild }; +} + +describe("RustLibrespotBackend.start", () => { + it("spawns librespot with the pipe/stdout arg set and the OAuth access token", async () => { + const h = makeHarness(); + await h.backend.start(); + expect(h.spawn).toHaveBeenCalledWith( + "/bin/librespot", + [ + "--name", "Test Bot", + "--backend", "pipe", + "--bitrate", "320", + "--format", "S16", + "--cache", "/tmp/cache", + "--device-type", "speaker", + "--access-token", "tok-123", + ], + expect.anything(), + ); + // NO --device (=> stdout) and NO --passthrough (=> decoded PCM, not Ogg). + const args = h.spawn.mock.calls.find((c) => String(c[0]).includes("librespot"))![1] as string[]; + expect(args).not.toContain("--device"); + expect(args).not.toContain("--passthrough"); + h.backend.stop(); + }); + + it("spawns ffmpeg (reader) before librespot (writer) with the exact 44100->48000 s16le args", async () => { + const h = makeHarness(); + await h.backend.start(); + const ffmpegArgs = h.spawn.mock.calls.find((c) => c[0] === "ffmpeg")![1] as string[]; + expect(ffmpegArgs).toEqual([ + "-f", "s16le", "-ar", "44100", "-ac", "2", "-i", "pipe:0", + "-f", "s16le", "-ar", "48000", "-ac", "2", "-acodec", "pcm_s16le", "pipe:1", + ]); + const ffmpegIdx = h.calls.indexOf("spawn:ffmpeg"); + const librespotIdx = h.calls.indexOf("spawn:librespot"); + expect(ffmpegIdx).toBeGreaterThanOrEqual(0); + expect(librespotIdx).toBeGreaterThan(ffmpegIdx); + h.backend.stop(); + }); + + it("getPcmStream() returns the ffmpeg stdout Readable", async () => { + const h = makeHarness(); + await h.backend.start(); + expect(h.backend.getPcmStream()).toBe(h.ffmpegChild.stdout); + h.backend.stop(); + }); + + it("emits 'ready' and reports isReady() true once our device appears in getDevices()", async () => { + const h = makeHarness(); + const ready = vi.fn(); + h.backend.on("ready", ready); + await h.backend.start(); + expect(h.connect.getDevices).toHaveBeenCalled(); + expect(ready).toHaveBeenCalledTimes(1); + expect(h.backend.isReady()).toBe(true); + h.backend.stop(); + }); + + it("keeps polling getDevices() until the device name appears", async () => { + const h = makeHarness(); + h.connect.getDevices + .mockResolvedValueOnce([]) + .mockResolvedValueOnce([{ id: "other", name: "Someone else", is_active: true }]) + .mockResolvedValue([{ id: "dev1", name: "Test Bot", is_active: false }]); + await h.backend.start(); + expect(h.connect.getDevices).toHaveBeenCalledTimes(3); + expect(h.backend.isReady()).toBe(true); + h.backend.stop(); + }); + + it("throws (and does not spawn) when the OAuth token is null", async () => { + const oauth = makeOAuth(); + oauth.getAccessToken.mockResolvedValue(null); + const h = makeHarness({ oauth }); + await expect(h.backend.start()).rejects.toThrow(/authorized|token/i); + expect(h.spawn).not.toHaveBeenCalled(); + }); +}); + +describe("RustLibrespotBackend transport delegation (Connect API)", () => { + it("playTrack resolves the device then transfer(false) then play(uri)", async () => { + const h = makeHarness(); + await h.backend.playTrack("spotify:track:go"); + expect(h.connect.findDeviceByName).toHaveBeenCalledWith("Test Bot"); + expect(h.connect.transfer).toHaveBeenCalledWith("dev1", false); + expect(h.connect.play).toHaveBeenCalledWith("dev1", "spotify:track:go"); + // ordering: transfer before play + expect(h.connect.transfer.mock.invocationCallOrder[0]) + .toBeLessThan(h.connect.play.mock.invocationCallOrder[0]); + }); + + it("playTrack throws when the device cannot be found", async () => { + const h = makeHarness(); + h.connect.findDeviceByName.mockResolvedValue(null); + await expect(h.backend.playTrack("spotify:track:x")).rejects.toThrow(/device/i); + }); + + it("pause/resume/seek delegate to the Connect API and seek updates position", async () => { + const h = makeHarness(); + await h.backend.pause(); + await h.backend.resume(); + await h.backend.seek(5000); + expect(h.connect.pause).toHaveBeenCalled(); + expect(h.connect.resume).toHaveBeenCalled(); + expect(h.connect.seek).toHaveBeenCalledWith(5000); + expect(h.backend.getPositionMs()).toBe(5000); + }); +}); + +describe("RustLibrespotBackend track-end poll loop", () => { + it("emits trackEnded when progress reaches the end-of-track window", async () => { + const h = makeHarness(); + const ended = vi.fn(); + const meta = vi.fn(); + h.backend.on("trackEnded", ended); + h.backend.on("metadata", meta); + h.connect.getPlaybackState + .mockResolvedValueOnce({ isPlaying: true, progressMs: 1000, trackUri: "spotify:track:A", durationMs: 200000 }) + .mockResolvedValueOnce({ isPlaying: true, progressMs: 199000, trackUri: "spotify:track:A", durationMs: 200000 }); + await (h.backend as any).pollState(); + expect(meta).toHaveBeenCalledWith(expect.objectContaining({ uri: "spotify:track:A", durationMs: 200000 })); + expect(h.backend.getPositionMs()).toBe(1000); + expect(ended).not.toHaveBeenCalled(); + await (h.backend as any).pollState(); + expect(ended).toHaveBeenCalledWith({ uri: "spotify:track:A", reason: "ended" }); + }); + + it("emits trackEnded once when playback stops (!isPlaying) after having played", async () => { + const h = makeHarness(); + const ended = vi.fn(); + h.backend.on("trackEnded", ended); + 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: 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(); // idempotent: no second emit for same track + expect(ended).toHaveBeenCalledTimes(1); + expect(ended).toHaveBeenCalledWith({ uri: "spotify:track:A", reason: "ended" }); + }); + + it("emits trackEnded when the track uri transitions to null after playing", async () => { + const h = makeHarness(); + const ended = vi.fn(); + h.backend.on("trackEnded", ended); + h.connect.getPlaybackState + .mockResolvedValueOnce({ isPlaying: true, progressMs: 1000, trackUri: "spotify:track:A", durationMs: 200000 }) + .mockResolvedValueOnce({ isPlaying: true, progressMs: 0, trackUri: null, durationMs: 0 }); + await (h.backend as any).pollState(); + await (h.backend as any).pollState(); + expect(ended).toHaveBeenCalledWith({ uri: "spotify:track:A", reason: "ended" }); + }); + + it("ignores a null playback state (no active device) without emitting", async () => { + const h = makeHarness(); + const ended = vi.fn(); + h.backend.on("trackEnded", ended); + h.connect.getPlaybackState.mockResolvedValue(null); + await (h.backend as any).pollState(); + expect(ended).not.toHaveBeenCalled(); + }); +}); + +describe("RustLibrespotBackend.stop", () => { + it("kills librespot + ffmpeg, clears ready, and is idempotent", async () => { + const h = makeHarness(); + await h.backend.start(); + h.backend.stop(); + h.backend.stop(); // second call must not throw + expect(h.librespotChild.kill).toHaveBeenCalled(); + expect(h.ffmpegChild.kill).toHaveBeenCalled(); + expect(h.backend.isReady()).toBe(false); + }); +}); + +describe("RustLibrespotBackend.start failure cleanup", () => { + it("tears down librespot + ffmpeg when the device never appears", async () => { + const h = makeHarness(); + h.connect.getDevices.mockResolvedValue([]); // device never shows up -> waitForDevice times out + await expect(h.backend.start()).rejects.toThrow(/did not appear/i); + expect(h.librespotChild.kill).toHaveBeenCalled(); + expect(h.ffmpegChild.kill).toHaveBeenCalled(); + expect(h.backend.isReady()).toBe(false); + }); +}); + +describe("RustLibrespotBackend child-process error handling", () => { + it("swallows+logs a child 'error' when no backend 'error' listener is attached", async () => { + const h = makeHarness(); + await h.backend.start(); + expect(h.backend.listenerCount("error")).toBe(0); + expect(() => h.librespotChild.emit("error", new Error("boom"))).not.toThrow(); + expect(() => h.ffmpegChild.emit("error", new Error("boom"))).not.toThrow(); + h.backend.stop(); + }); + + it("re-emits a child 'error' to an attached backend 'error' listener", async () => { + const h = makeHarness(); + await h.backend.start(); + const onErr = vi.fn(); + h.backend.on("error", onErr); + const err = new Error("ffmpeg boom"); + h.ffmpegChild.emit("error", err); + expect(onErr).toHaveBeenCalledWith(err); + h.backend.stop(); + }); +}); diff --git a/src/music/spotify/rust-librespot.ts b/src/music/spotify/rust-librespot.ts new file mode 100644 index 0000000..abbb544 --- /dev/null +++ b/src/music/spotify/rust-librespot.ts @@ -0,0 +1,332 @@ +import { EventEmitter } from "node:events"; +import type { Readable } from "node:stream"; +import type { ChildProcess } from "node:child_process"; +import { spawn as realSpawn } from "node:child_process"; +import { mkdirSync as realMkdirSync } from "node:fs"; +import type { Logger } from "pino"; +import type { + SpotifyAudioBackend, + SpotifyTrackEndedEvent, + SpotifyNowPlaying, +} from "./backend.js"; +import { findLibrespot } from "./binary.js"; +import { SpotifyConnectApi } from "./connect-api.js"; +import type { PlaybackState, SpotifyDevice } from "./connect-api.js"; +import type { SpotifyOAuth } from "./spotify-oauth.js"; +import { getFfmpegCommand } from "../../audio/player.js"; + +export interface RustLibrespotBackendOptions { + deviceName: string; + bitrate: number; + cacheDir: string; + oauth: SpotifyOAuth; + connect?: SpotifyConnectApi; + logger: Logger; + deps?: RustLibrespotBackendDeps; +} + +/** Injectable seams so the whole lifecycle is testable without a real binary/network. */ +export interface RustLibrespotBackendDeps { + spawn?: typeof realSpawn; + mkdirSync?: typeof realMkdirSync; + findBinary?: () => string; + /** + * C1: override the ffmpeg command. Production resolves it via + * getFfmpegCommand() (bundled ffmpeg-static fallback when `ffmpeg` isn't on + * PATH); tests pin it to "ffmpeg" for stable arg assertions. + */ + ffmpegCommand?: string; + sleep?: (ms: number) => Promise; + readyPollIntervalMs?: number; + readyTimeoutMs?: number; + statePollIntervalMs?: number; +} + +const DEFAULT_READY_POLL_MS = 500; +const DEFAULT_READY_TIMEOUT_MS = 20_000; +const DEFAULT_STATE_POLL_MS = 2_000; +/** How close to the end (ms) counts as "track finished" when polling player state. */ +const END_OF_TRACK_WINDOW_MS = 1_500; + +const defaultSleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); + +export class RustLibrespotBackend extends EventEmitter implements SpotifyAudioBackend { + private readonly opts: RustLibrespotBackendOptions; + private readonly log: Logger; + private readonly deps: RustLibrespotBackendDeps; + private readonly oauth: SpotifyOAuth; + private readonly connect: SpotifyConnectApi; + + private proc: ChildProcess | null = null; + private ffmpeg: ChildProcess | null = null; + private pollTimer: ReturnType | null = null; + private ready = false; + private positionMs = 0; + + // track-end poll state machine + private currentUri: string | null = null; + private hasPlayed = false; + private endedForCurrent = false; + + constructor(o: RustLibrespotBackendOptions) { + super(); + this.opts = o; + this.log = o.logger; + this.deps = o.deps ?? {}; + this.oauth = o.oauth; + // The Connect API shares the backend's OAuth token source. Reuse the + // injected instance in tests; otherwise build one over oauth.getAccessToken(). + this.connect = o.connect ?? new SpotifyConnectApi(() => this.oauth.getAccessToken()); + } + + async start(): Promise { + const spawn = this.deps.spawn ?? realSpawn; + const mkdirSync = this.deps.mkdirSync ?? realMkdirSync; + const findBinary = this.deps.findBinary ?? findLibrespot; + // C1: resolve ffmpeg via getFfmpegCommand() unless injected for tests. + const ffmpegCommand = this.deps.ffmpegCommand ?? getFfmpegCommand(); + + // A valid USER control token is required before we spawn anything. + const token = await this.oauth.getAccessToken(); + if (!token) { + throw new Error("Spotify not authorized (no access token) — sign in first"); + } + + mkdirSync(this.opts.cacheDir, { recursive: true }); + + // 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 { + // 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( + ffmpegCommand, + [ + "-f", "s16le", "-ar", "44100", "-ac", "2", "-i", "pipe:0", + "-f", "s16le", "-ar", "48000", "-ac", "2", "-acodec", "pcm_s16le", "pipe:1", + ], + { stdio: ["pipe", "pipe", "pipe"] }, + ); + this.ffmpeg.stderr?.on("data", (b: Buffer) => + this.log.debug({ ffmpeg: b.toString().trim() }, "ffmpeg"), + ); + this.ffmpeg.on("error", (err) => this.emitError(err)); + + // 2. Spawn librespot: --backend pipe with NO --device => raw s16le/44100/2 + // on stdout, NO --passthrough (that would emit raw Ogg). --access-token + // authenticates it as a Connect device controllable via the Web API. + const bin = findBinary(); + this.proc = spawn( + bin, + [ + "--name", this.opts.deviceName, + "--backend", "pipe", + "--bitrate", String(this.opts.bitrate), + "--format", "S16", + "--cache", this.opts.cacheDir, + "--device-type", "speaker", + "--access-token", token, + ], + { stdio: ["ignore", "pipe", "pipe"] }, + ); + // 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); + } + this.proc.stderr?.on("data", (b: Buffer) => + this.log.info({ librespot: b.toString().trim() }, "librespot"), + ); + this.proc.on("error", (err) => this.emitError(err)); + this.proc.on("exit", (code, signal) => { + this.ready = false; + this.log.warn({ code, signal }, "librespot exited"); + }); + + // 3. Poll the Connect device list until our device registers. + await this.waitForDevice(); + + // 4. Begin the player-state poll loop (track-end / position / metadata). + this.startPollLoop(); + + this.ready = true; + this.emit("ready"); + } catch (e) { + this.stop(); + throw e; + } + } + + /** + * Re-emit a child "error" only when a consumer is listening; Node throws on an + * unhandled "error" event, so with no listener we log via the injected logger. + */ + private emitError(err: unknown): void { + if (this.listenerCount("error") > 0) { + this.emit("error", err); + } else { + this.log.error({ err }, "rust-librespot backend error (no listener)"); + } + } + + private async waitForDevice(): Promise { + const sleep = this.deps.sleep ?? defaultSleep; + const interval = this.deps.readyPollIntervalMs ?? DEFAULT_READY_POLL_MS; + const timeout = this.deps.readyTimeoutMs ?? DEFAULT_READY_TIMEOUT_MS; + const deadline = Date.now() + timeout; + while (Date.now() < deadline) { + let devices: SpotifyDevice[] = []; + try { + devices = await this.connect.getDevices(); + } catch (err) { + this.log.debug({ err }, "getDevices failed during readiness poll"); + } + if (devices.some((d) => d.name === this.opts.deviceName)) return; + await sleep(interval); + } + throw new Error(`librespot device "${this.opts.deviceName}" did not appear within timeout`); + } + + private startPollLoop(): void { + const interval = this.deps.statePollIntervalMs ?? DEFAULT_STATE_POLL_MS; + this.pollTimer = setInterval(() => { + void this.pollState(); + }, interval); + // Don't keep the event loop / test process alive on account of the poll timer. + this.pollTimer.unref?.(); + } + + /** One player-state poll iteration: updates position/metadata and detects track end. */ + private async pollState(): Promise { + let state: PlaybackState | null; + try { + state = await this.connect.getPlaybackState(); + } catch (err) { + this.log.debug({ err }, "getPlaybackState failed"); + return; + } + + if (!state) { + // C3.5: a 204 / no-active-device response. AFTER our own track has been + // seen playing, librespot going idle means the track ended — emit once so + // the queue advances instead of stalling. BEFORE any play (hasPlayed + // false), a null state is just "nothing active yet" and is ignored. + if (this.hasPlayed && this.currentUri && !this.endedForCurrent) { + this.endedForCurrent = true; + const endedUri = this.currentUri; + this.currentUri = null; + const e: SpotifyTrackEndedEvent = { uri: endedUri, reason: "ended" }; + this.emit("trackEnded", e); + } + return; + } + + this.positionMs = state.progressMs; + + // Track change -> reset the end-detection state and surface best-effort metadata. + if (state.trackUri && state.trackUri !== this.currentUri) { + this.currentUri = state.trackUri; + this.hasPlayed = false; + this.endedForCurrent = false; + const np: SpotifyNowPlaying = { + uri: state.trackUri, + name: "", + artist: "", + album: "", + coverUrl: "", + durationMs: state.durationMs, + }; + this.emit("metadata", np); + } + + if (state.isPlaying) this.hasPlayed = true; + if (!this.currentUri || this.endedForCurrent) return; + + // C3.4: EVERY end condition is gated on hasPlayed so no end can fire until + // 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". + const finishedByProgress = + this.hasPlayed && + state.durationMs > 0 && + state.progressMs >= state.durationMs - END_OF_TRACK_WINDOW_MS; + const finishedByStop = this.hasPlayed && !state.isPlaying; + const finishedByNull = this.hasPlayed && state.trackUri === null; + + if (finishedByProgress || finishedByStop || finishedByNull) { + this.endedForCurrent = true; // latch: emit at most once per track + const endedUri = this.currentUri; + this.currentUri = null; + const e: SpotifyTrackEndedEvent = { uri: endedUri, reason: "ended" }; + this.emit("trackEnded", e); + } + } + + isReady(): boolean { + return this.ready; + } + + async playTrack(uri: string): Promise { + const deviceId = await this.connect.findDeviceByName(this.opts.deviceName); + if (!deviceId) throw new Error(`Connect device "${this.opts.deviceName}" not found`); + // Reset the track-end state machine for the new track: clear the once-only + // latch and drop hasPlayed so no end can fire until a poll re-confirms this + // uri playing. currentUri is cleared so the next poll re-detects the track + // (fresh metadata) rather than treating it as unchanged. + this.currentUri = null; + this.hasPlayed = false; + this.endedForCurrent = 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. + await this.connect.transfer(deviceId, false); + await this.connect.play(deviceId, uri); + } + + async pause(): Promise { + await this.connect.pause(); + } + + async resume(): Promise { + await this.connect.resume(); + } + + async seek(ms: number): Promise { + await this.connect.seek(ms); + this.positionMs = ms; + } + + getPcmStream(): Readable { + const out = this.ffmpeg?.stdout; + if (!out) throw new Error("PCM stream unavailable (rust-librespot backend not started)"); + return out; + } + + getPositionMs(): number { + return this.positionMs; + } + + stop(): void { + this.ready = false; + // Clear the state poll interval FIRST so no poll fires mid-teardown. + if (this.pollTimer) { + clearInterval(this.pollTimer); + this.pollTimer = null; + } + if (this.proc) { + try { + this.proc.kill(); + } catch { + /* ignore */ + } + this.proc = null; + } + if (this.ffmpeg) { + try { + this.ffmpeg.kill(); + } catch { + /* ignore */ + } + this.ffmpeg = null; + } + } +}