diff --git a/src/audio/player.ts b/src/audio/player.ts index ffc4a9c..36d0231 100644 --- a/src/audio/player.ts +++ b/src/audio/player.ts @@ -1,10 +1,19 @@ import { spawn, type ChildProcess } from "node:child_process"; import { EventEmitter } from "node:events"; import { createOpusEncoder, PCM_FRAME_BYTES, type Encoder } from "./encoder.js"; -import type pino from "pino"; +import type { Logger } from "../logger.js"; + +export interface PlayerEvents { + frame: (opusFrame: Buffer) => void; + trackEnd: () => void; + error: (err: Error) => void; +} export type PlayerState = "idle" | "playing" | "paused"; +// ~5 seconds of 48kHz stereo 16-bit audio +const PCM_HIGH_WATER_MARK = PCM_FRAME_BYTES * 250; + export class AudioPlayer extends EventEmitter { private ffmpeg: ChildProcess | null = null; private encoder: Encoder; @@ -12,9 +21,9 @@ export class AudioPlayer extends EventEmitter { private volume = 75; // 0-100 private frameTimer: ReturnType | null = null; private pcmBuffer: Buffer = Buffer.alloc(0); - private logger: pino.Logger; + private logger: Logger; - constructor(logger: pino.Logger) { + constructor(logger: Logger) { super(); this.encoder = createOpusEncoder(); this.logger = logger; @@ -43,6 +52,10 @@ export class AudioPlayer extends EventEmitter { this.ffmpeg.stdout!.on("data", (chunk: Buffer) => { this.pcmBuffer = Buffer.concat([this.pcmBuffer, chunk]); + // Backpressure: pause FFmpeg if buffer is too large + if (this.pcmBuffer.length > PCM_HIGH_WATER_MARK && this.ffmpeg?.stdout) { + this.ffmpeg.stdout.pause(); + } }); this.ffmpeg.on("close", (code) => { @@ -73,6 +86,14 @@ export class AudioPlayer extends EventEmitter { const pcmFrame = this.pcmBuffer.subarray(0, PCM_FRAME_BYTES); this.pcmBuffer = this.pcmBuffer.subarray(PCM_FRAME_BYTES); + // Resume FFmpeg if buffer drained below threshold + if ( + this.pcmBuffer.length < PCM_HIGH_WATER_MARK / 2 && + this.ffmpeg?.stdout?.isPaused() + ) { + this.ffmpeg.stdout.resume(); + } + const opusFrame = this.encoder.encode(Buffer.from(pcmFrame)); this.emit("frame", opusFrame); } diff --git a/src/logger.ts b/src/logger.ts index d857d1a..5eebfc4 100644 --- a/src/logger.ts +++ b/src/logger.ts @@ -2,7 +2,9 @@ import pino from "pino"; import { mkdirSync } from "node:fs"; import { join } from "node:path"; -export function createLogger(logDir?: string): pino.Logger { +export type Logger = pino.Logger; + +export function createLogger(logDir?: string): Logger { if (!logDir) { return pino({ level: "info" }); } diff --git a/src/ts-protocol/client.ts b/src/ts-protocol/client.ts index 15e6c88..2e6dea2 100644 --- a/src/ts-protocol/client.ts +++ b/src/ts-protocol/client.ts @@ -7,7 +7,7 @@ import { exportIdentity, type TS3Identity, } from "./identity.js"; -import type pino from "pino"; +import type { Logger } from "../logger.js"; export interface TS3ClientOptions { host: string; @@ -33,9 +33,9 @@ export class TS3Client extends EventEmitter { private identity: TS3Identity; private clientId = 0; private keepAliveInterval: ReturnType | null = null; - private logger: pino.Logger; + private logger: Logger; - constructor(private options: TS3ClientOptions, logger: pino.Logger) { + constructor(private options: TS3ClientOptions, logger: Logger) { super(); this.logger = logger;