feat(voice): add voice transmit functionality for real-time audio streaming
- Add VoiceTransmitter class to handle PCM audio from backend/browser to Discord - Implement voice:transmit:start and voice:transmit:stop command handlers - Upsample 24kHz mono PCM to 48kHz stereo for Discord compatibility - Encode PCM to Opus and stream via browser-bridge player owner - Subscribe to Redis channel backend:voice:transmit for real-time PCM data Features: - Backend can send PCM audio (24kHz mono s16le base64) via Redis - Automatic upsampling and encoding to Discord-compatible format - Clean start/stop lifecycle with resource cleanup This completes the bidirectional voice streaming: - Listen: Discord → Backend (already working via voicePcmData broadcast) - Transmit: Backend → Discord (now implemented) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
60edb4576f
commit
f74252f594
@@ -3,6 +3,7 @@ import type { Client } from "discord.js-selfbot-v13";
|
||||
import Redis from "ioredis";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { discordPlayer } from "../voice-recording/player.js";
|
||||
import { voiceTransmitter } from "../voice-recording/transmitter.js";
|
||||
import type { VoiceController } from "../voice-recording/voiceController.js";
|
||||
|
||||
const logger = createChildLogger("command-handler");
|
||||
@@ -130,6 +131,12 @@ export class CommandHandler {
|
||||
case "voice:channels":
|
||||
reply = await this.handleVoiceChannels(cmd);
|
||||
break;
|
||||
case "voice:transmit:start":
|
||||
reply = await this.handleVoiceTransmitStart(cmd);
|
||||
break;
|
||||
case "voice:transmit:stop":
|
||||
reply = await this.handleVoiceTransmitStop(cmd);
|
||||
break;
|
||||
case "guilds:list":
|
||||
reply = await this.handleListGuilds(cmd);
|
||||
break;
|
||||
@@ -370,6 +377,63 @@ export class CommandHandler {
|
||||
};
|
||||
}
|
||||
|
||||
private async handleVoiceTransmitStart(cmd: BackendCommand): Promise<CommandReply> {
|
||||
if (!discordPlayer.isConnected()) {
|
||||
return {
|
||||
id: cmd.id,
|
||||
success: false,
|
||||
data: null,
|
||||
error: "Not connected to voice channel",
|
||||
};
|
||||
}
|
||||
|
||||
try {
|
||||
// Create a new Redis connection for the transmitter
|
||||
const transmitRedis = new Redis(config.REDIS_URL);
|
||||
await voiceTransmitter.start(transmitRedis);
|
||||
|
||||
const status = voiceTransmitter.getStatus();
|
||||
logger.info({ status }, "Voice transmit started");
|
||||
|
||||
return {
|
||||
id: cmd.id,
|
||||
success: true,
|
||||
data: status,
|
||||
};
|
||||
} catch (err) {
|
||||
const message = err instanceof Error ? err.message : String(err);
|
||||
logger.error({ error: message }, "Failed to start voice transmit");
|
||||
return {
|
||||
id: cmd.id,
|
||||
success: false,
|
||||
data: null,
|
||||
error: message,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private async handleVoiceTransmitStop(cmd: BackendCommand): Promise<CommandReply> {
|
||||
try {
|
||||
await voiceTransmitter.stop();
|
||||
logger.info("Voice transmit stopped");
|
||||
|
||||
return {
|
||||
id: cmd.id,
|
||||
success: true,
|
||||
data: { status: "stopped" },
|
||||
};
|
||||
} catch (err) {
|
||||
const message = err instanceof Error ? err.message : String(err);
|
||||
logger.error({ error: message }, "Failed to stop voice transmit");
|
||||
return {
|
||||
id: cmd.id,
|
||||
success: false,
|
||||
data: null,
|
||||
error: message,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
// ---- Status publishing ----
|
||||
|
||||
private publishVoiceStatus(): void {
|
||||
|
||||
@@ -2,3 +2,4 @@ export { OpusDecoder } from "./recorder/decoder.js";
|
||||
export { SegmentManager } from "./recorder/segment.js";
|
||||
export { startRecording, stopRecording } from "./recorder.js";
|
||||
export { VoiceController } from "./voiceController.js";
|
||||
export { voiceTransmitter } from "./transmitter.js";
|
||||
|
||||
@@ -0,0 +1,168 @@
|
||||
import { PassThrough, Readable } from "node:stream";
|
||||
import { createChildLogger } from "@bete/shared/logger";
|
||||
import { StreamType } from "@discordjs/voice";
|
||||
import type Redis from "ioredis";
|
||||
import prism from "prism-media";
|
||||
import { discordPlayer } from "./player.js";
|
||||
|
||||
const logger = createChildLogger("transmitter");
|
||||
|
||||
/**
|
||||
* Handles real-time PCM audio transmission from backend/browser to Discord.
|
||||
* Receives 24kHz mono PCM data, upsamples to 48kHz stereo, encodes to Opus, and plays to Discord.
|
||||
*/
|
||||
export class VoiceTransmitter {
|
||||
private redisSub: Redis | null = null;
|
||||
private pcmStream: PassThrough | null = null;
|
||||
private opusEncoder: any | null = null;
|
||||
private isActive = false;
|
||||
private readonly TRANSMIT_CHANNEL = "backend:voice:transmit";
|
||||
|
||||
/**
|
||||
* Start listening for PCM audio data from Redis and stream to Discord
|
||||
*/
|
||||
async start(redis: Redis): Promise<void> {
|
||||
if (this.isActive) {
|
||||
logger.warn("Voice transmitter already active");
|
||||
return;
|
||||
}
|
||||
|
||||
this.redisSub = redis;
|
||||
this.isActive = true;
|
||||
|
||||
// Create PCM input stream
|
||||
this.pcmStream = new PassThrough();
|
||||
|
||||
// Upsample 24kHz mono → 48kHz stereo
|
||||
const upsampledStream = this.upsampleTo48kStereo(this.pcmStream);
|
||||
|
||||
// Encode to Opus
|
||||
this.opusEncoder = new prism.opus.Encoder({
|
||||
rate: 48000,
|
||||
channels: 2,
|
||||
frameSize: 960,
|
||||
});
|
||||
|
||||
const opusStream = upsampledStream.pipe(this.opusEncoder);
|
||||
|
||||
// Play to Discord
|
||||
discordPlayer.playStream(opusStream, "browser-bridge", {
|
||||
inputType: StreamType.Opus,
|
||||
inlineVolume: false,
|
||||
});
|
||||
|
||||
// Subscribe to Redis channel for PCM data
|
||||
await this.redisSub.subscribe(this.TRANSMIT_CHANNEL);
|
||||
|
||||
this.redisSub.on("message", (channel, message) => {
|
||||
if (channel !== this.TRANSMIT_CHANNEL || !this.pcmStream) return;
|
||||
|
||||
try {
|
||||
const data = JSON.parse(message);
|
||||
if (data.type === "pcm" && data.buffer) {
|
||||
const pcmBuffer = Buffer.from(data.buffer, "base64");
|
||||
this.pcmStream.write(pcmBuffer);
|
||||
}
|
||||
} catch (err) {
|
||||
logger.error({ error: err }, "Failed to process PCM data");
|
||||
}
|
||||
});
|
||||
|
||||
logger.info("Voice transmitter started");
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop transmitting and clean up resources
|
||||
*/
|
||||
async stop(): Promise<void> {
|
||||
if (!this.isActive) return;
|
||||
|
||||
this.isActive = false;
|
||||
|
||||
if (this.pcmStream) {
|
||||
this.pcmStream.end();
|
||||
this.pcmStream = null;
|
||||
}
|
||||
|
||||
if (this.opusEncoder) {
|
||||
this.opusEncoder.destroy();
|
||||
this.opusEncoder = null;
|
||||
}
|
||||
|
||||
if (this.redisSub) {
|
||||
await this.redisSub.unsubscribe(this.TRANSMIT_CHANNEL);
|
||||
this.redisSub = null;
|
||||
}
|
||||
|
||||
discordPlayer.stop("browser-bridge");
|
||||
|
||||
logger.info("Voice transmitter stopped");
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if transmitter is currently active
|
||||
*/
|
||||
getStatus(): { active: boolean; channel: string } {
|
||||
return {
|
||||
active: this.isActive,
|
||||
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();
|
||||
Reference in New Issue
Block a user