From f74252f594f725fdb61d8e70a8800bee9f1f6e5d Mon Sep 17 00:00:00 2001 From: MythEclipse Date: Mon, 8 Jun 2026 20:48:16 +0700 Subject: [PATCH] feat(voice): add voice transmit functionality for real-time audio streaming MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- .../modules/command-handler/commandHandler.ts | 64 +++++++ .../src/modules/voice-recording/index.ts | 1 + .../modules/voice-recording/transmitter.ts | 168 ++++++++++++++++++ 3 files changed, 233 insertions(+) create mode 100644 services/discord-gateway/src/modules/voice-recording/transmitter.ts diff --git a/services/discord-gateway/src/modules/command-handler/commandHandler.ts b/services/discord-gateway/src/modules/command-handler/commandHandler.ts index c43f0d0..b7208c1 100644 --- a/services/discord-gateway/src/modules/command-handler/commandHandler.ts +++ b/services/discord-gateway/src/modules/command-handler/commandHandler.ts @@ -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 { + 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 { + 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 { diff --git a/services/discord-gateway/src/modules/voice-recording/index.ts b/services/discord-gateway/src/modules/voice-recording/index.ts index f435647..39bb99d 100644 --- a/services/discord-gateway/src/modules/voice-recording/index.ts +++ b/services/discord-gateway/src/modules/voice-recording/index.ts @@ -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"; diff --git a/services/discord-gateway/src/modules/voice-recording/transmitter.ts b/services/discord-gateway/src/modules/voice-recording/transmitter.ts new file mode 100644 index 0000000..d7cbd10 --- /dev/null +++ b/services/discord-gateway/src/modules/voice-recording/transmitter.ts @@ -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 { + 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 { + 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();