fix(voice): use FFmpeg for OggOpus encoding in transmitter
- Replace prism Opus.Encoder + OggLogicalBitstream with FFmpeg - FFmpeg handles upsampling, encoding, and OGG container in one process - Input: raw PCM 24kHz mono s16le via stdin - Output: OggOpus via stdout (StreamType.OggOpus) - FFmpeg arguments optimized for real-time low-delay streaming Previous approach failed because: 1. Raw Opus packets without OGG wrapper don't work with @discordjs/voice 2. prism's OggLogicalBitstream has CRC bug with node-crc native bindings 3. Manual upsampling was error-prone Pipeline: Browser Mic → base64 PCM → Redis → FFmpeg → OggOpus → Discord Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
25eb3f6353
commit
942605661a
@@ -1,20 +1,23 @@
|
|||||||
import { PassThrough, Readable } from "node:stream";
|
import { PassThrough } from "node:stream";
|
||||||
|
import { spawn } from "node:child_process";
|
||||||
import { createChildLogger } from "@bete/shared/logger";
|
import { createChildLogger } from "@bete/shared/logger";
|
||||||
import { StreamType } from "@discordjs/voice";
|
import { StreamType } from "@discordjs/voice";
|
||||||
import type Redis from "ioredis";
|
import type Redis from "ioredis";
|
||||||
import prism from "prism-media";
|
|
||||||
import { discordPlayer } from "./player.js";
|
import { discordPlayer } from "./player.js";
|
||||||
|
|
||||||
const logger = createChildLogger("transmitter");
|
const logger = createChildLogger("transmitter");
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Handles real-time PCM audio transmission from backend/browser to Discord.
|
* Handles real-time PCM audio transmission from browser to Discord voice channel.
|
||||||
* Receives 24kHz mono PCM data, upsamples to 48kHz stereo, encodes to Opus, and plays to Discord.
|
*
|
||||||
|
* Pipeline:
|
||||||
|
* Browser Mic → base64 PCM (24kHz mono s16le) → Redis →
|
||||||
|
* FFmpeg (encode to OggOpus) → discordPlayer (StreamType.OggOpus) → Discord Voice
|
||||||
*/
|
*/
|
||||||
export class VoiceTransmitter {
|
export class VoiceTransmitter {
|
||||||
private redisSub: Redis | null = null;
|
private redisSub: Redis | null = null;
|
||||||
private pcmStream: PassThrough | null = null;
|
private pcmStream: PassThrough | null = null;
|
||||||
private opusEncoder: any | null = null;
|
private ffmpegProcess: ReturnType<typeof spawn> | null = null;
|
||||||
private isActive = false;
|
private isActive = false;
|
||||||
private readonly TRANSMIT_CHANNEL = "backend:voice:transmit";
|
private readonly TRANSMIT_CHANNEL = "backend:voice:transmit";
|
||||||
|
|
||||||
@@ -33,26 +36,62 @@ export class VoiceTransmitter {
|
|||||||
// Create PCM input stream
|
// Create PCM input stream
|
||||||
this.pcmStream = new PassThrough();
|
this.pcmStream = new PassThrough();
|
||||||
|
|
||||||
// Upsample 24kHz mono → 48kHz stereo
|
// Spawn FFmpeg to encode 24kHz mono PCM → OggOpus
|
||||||
const upsampledStream = this.upsampleTo48kStereo(this.pcmStream);
|
// Input: 24kHz mono s16le (raw PCM)
|
||||||
|
// Output: OGG container with Opus audio
|
||||||
// Encode to Opus — Encoder outputs raw Opus packets (no OGG wrapper needed)
|
this.ffmpegProcess = spawn("ffmpeg", [
|
||||||
this.opusEncoder = new prism.opus.Encoder({
|
"-f", "s16le", // Input format: signed 16-bit little-endian
|
||||||
rate: 48000,
|
"-ar", "24000", // Input sample rate: 24kHz
|
||||||
channels: 2,
|
"-ac", "1", // Input channels: mono
|
||||||
frameSize: 960,
|
"-i", "pipe:0", // Read from stdin
|
||||||
|
"-f", "ogg", // Output format: OGG
|
||||||
|
"-c:a", "libopus", // Codec: Opus
|
||||||
|
"-b:a", "96k", // Bitrate: 96kbps
|
||||||
|
"-ar", "48000", // Output sample rate: 48kHz
|
||||||
|
"-ac", "2", // Output channels: stereo
|
||||||
|
"-application", "lowdelay", // Low delay mode for real-time
|
||||||
|
"-frame_duration", "20", // 20ms frames
|
||||||
|
"-packet_loss", "0", // No packet loss expected
|
||||||
|
"pipe:1", // Write to stdout
|
||||||
|
], {
|
||||||
|
stdio: ["pipe", "pipe", "pipe"],
|
||||||
});
|
});
|
||||||
|
|
||||||
const opusStream = upsampledStream.pipe(this.opusEncoder);
|
// Pipe PCM data to FFmpeg stdin
|
||||||
|
if (this.ffmpegProcess.stdin) {
|
||||||
|
this.pcmStream.pipe(this.ffmpegProcess.stdin);
|
||||||
|
}
|
||||||
|
|
||||||
// Play to Discord with raw Opus format
|
// Log FFmpeg stderr for debugging
|
||||||
discordPlayer.playStream(opusStream, "browser-bridge", {
|
const stderrChunks: Buffer[] = [];
|
||||||
inputType: StreamType.Opus,
|
this.ffmpegProcess.stderr?.on("data", (chunk: Buffer) => {
|
||||||
|
stderrChunks.push(chunk);
|
||||||
|
});
|
||||||
|
|
||||||
|
this.ffmpegProcess.on("error", (err) => {
|
||||||
|
logger.error({ error: err.message }, "FFmpeg process error");
|
||||||
|
});
|
||||||
|
|
||||||
|
this.ffmpegProcess.on("exit", (code) => {
|
||||||
|
if (code !== 0) {
|
||||||
|
const stderr = Buffer.concat(stderrChunks).toString();
|
||||||
|
logger.error({ code, stderr: stderr.slice(-500) }, "FFmpeg exited with error");
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// Play FFmpeg stdout (OggOpus) to Discord
|
||||||
|
if (this.ffmpegProcess.stdout) {
|
||||||
|
discordPlayer.playStream(this.ffmpegProcess.stdout, "browser-bridge", {
|
||||||
|
inputType: StreamType.OggOpus,
|
||||||
inlineVolume: true,
|
inlineVolume: true,
|
||||||
});
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
logger.info("Voice transmitter pipeline ready (PCM → FFmpeg → OggOpus → Discord)");
|
||||||
|
|
||||||
// Subscribe to Redis channel for PCM data
|
// Subscribe to Redis channel for PCM data
|
||||||
await this.redisSub.subscribe(this.TRANSMIT_CHANNEL);
|
await this.redisSub.subscribe(this.TRANSMIT_CHANNEL);
|
||||||
|
logger.info({ channel: this.TRANSMIT_CHANNEL }, "Subscribed to transmit channel");
|
||||||
|
|
||||||
this.redisSub.on("message", (channel, message) => {
|
this.redisSub.on("message", (channel, message) => {
|
||||||
if (channel !== this.TRANSMIT_CHANNEL || !this.pcmStream) return;
|
if (channel !== this.TRANSMIT_CHANNEL || !this.pcmStream) return;
|
||||||
@@ -61,6 +100,7 @@ export class VoiceTransmitter {
|
|||||||
const data = JSON.parse(message);
|
const data = JSON.parse(message);
|
||||||
if (data.type === "pcm" && data.buffer) {
|
if (data.type === "pcm" && data.buffer) {
|
||||||
const pcmBuffer = Buffer.from(data.buffer, "base64");
|
const pcmBuffer = Buffer.from(data.buffer, "base64");
|
||||||
|
logger.debug({ bytes: pcmBuffer.length }, "Received PCM chunk");
|
||||||
this.pcmStream.write(pcmBuffer);
|
this.pcmStream.write(pcmBuffer);
|
||||||
}
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
@@ -84,9 +124,9 @@ export class VoiceTransmitter {
|
|||||||
this.pcmStream = null;
|
this.pcmStream = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (this.opusEncoder) {
|
if (this.ffmpegProcess) {
|
||||||
this.opusEncoder.destroy();
|
this.ffmpegProcess.kill("SIGTERM");
|
||||||
this.opusEncoder = null;
|
this.ffmpegProcess = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (this.redisSub) {
|
if (this.redisSub) {
|
||||||
@@ -95,7 +135,6 @@ export class VoiceTransmitter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
discordPlayer.stop("browser-bridge");
|
discordPlayer.stop("browser-bridge");
|
||||||
|
|
||||||
logger.info("Voice transmitter stopped");
|
logger.info("Voice transmitter stopped");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -108,61 +147,6 @@ export class VoiceTransmitter {
|
|||||||
channel: this.TRANSMIT_CHANNEL,
|
channel: this.TRANSMIT_CHANNEL,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* Upsample 24kHz mono PCM to 48kHz stereo
|
|
||||||
* Input: 24kHz mono s16le (2 bytes per sample)
|
|
||||||
* Output: 48kHz stereo s16le (4 bytes per sample)
|
|
||||||
*/
|
|
||||||
private upsampleTo48kStereo(input: Readable): Readable {
|
|
||||||
const output = new PassThrough();
|
|
||||||
|
|
||||||
input.on("data", (chunk: Buffer) => {
|
|
||||||
// 24kHz mono → 48kHz stereo means we need to:
|
|
||||||
// 1. Duplicate each sample (mono → stereo)
|
|
||||||
// 2. Interpolate samples (24kHz → 48kHz)
|
|
||||||
|
|
||||||
const inputSamples = chunk.length / 2; // 16-bit samples
|
|
||||||
const outputBuffer = Buffer.alloc(inputSamples * 4 * 2); // 2x rate, 2x channels
|
|
||||||
|
|
||||||
for (let i = 0; i < inputSamples; i++) {
|
|
||||||
const sample = chunk.readInt16LE(i * 2);
|
|
||||||
|
|
||||||
// Write to output at 2x rate with simple duplication
|
|
||||||
// Sample i → output[i*2] and output[i*2+1]
|
|
||||||
const outIdx = i * 2;
|
|
||||||
|
|
||||||
// Left channel
|
|
||||||
outputBuffer.writeInt16LE(sample, outIdx * 4);
|
|
||||||
// Right channel
|
|
||||||
outputBuffer.writeInt16LE(sample, outIdx * 4 + 2);
|
|
||||||
|
|
||||||
// Interpolated sample (simple average for smoothing)
|
|
||||||
if (i < inputSamples - 1) {
|
|
||||||
const nextSample = chunk.readInt16LE((i + 1) * 2);
|
|
||||||
const interpolated = Math.floor((sample + nextSample) / 2);
|
|
||||||
|
|
||||||
// Left channel
|
|
||||||
outputBuffer.writeInt16LE(interpolated, (outIdx + 1) * 4);
|
|
||||||
// Right channel
|
|
||||||
outputBuffer.writeInt16LE(interpolated, (outIdx + 1) * 4 + 2);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
output.write(outputBuffer);
|
|
||||||
});
|
|
||||||
|
|
||||||
input.on("end", () => {
|
|
||||||
output.end();
|
|
||||||
});
|
|
||||||
|
|
||||||
input.on("error", (err) => {
|
|
||||||
logger.error({ error: err }, "Upsample input stream error");
|
|
||||||
output.destroy(err);
|
|
||||||
});
|
|
||||||
|
|
||||||
return output;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export const voiceTransmitter = new VoiceTransmitter();
|
export const voiceTransmitter = new VoiceTransmitter();
|
||||||
|
|||||||
Reference in New Issue
Block a user