mirror of
https://github.com/ZHANGTIANYAO1/teamspeak-music-bot.git
synced 2026-10-02 04:52:50 +08:00
feat(spotify): add RustLibrespotBackend (stdout-pipe -> ffmpeg, Connect-API control)
Implements the Stage-2 SpotifyAudioBackend over Rust librespot: spawns librespot with --backend pipe (no --device => s16le/44100/2 on stdout, no --passthrough), pipes stdout -> ffmpeg (44100->48000 s16le), waits for the Connect device to register before emitting "ready", and controls playback (transfer/play/pause/resume/seek) plus track-end/position/metadata via a polled SpotifyConnectApi. child_process/connect/oauth/ffmpeg injected for fully mocked, no-network unit tests (Windows-targeted; not e2e without Premium). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
540bf8c032
commit
35bdd2a168
2 files changed
+615
No files matched your search
@@ -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();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -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<void>;
|
||||||
|
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<void>((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<typeof setInterval> | 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<void> {
|
||||||
|
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<void> {
|
||||||
|
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<void> {
|
||||||
|
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<void> {
|
||||||
|
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<void> {
|
||||||
|
await this.connect.pause();
|
||||||
|
}
|
||||||
|
|
||||||
|
async resume(): Promise<void> {
|
||||||
|
await this.connect.resume();
|
||||||
|
}
|
||||||
|
|
||||||
|
async seek(ms: number): Promise<void> {
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in new issue
Block a user