From 3965ea79bceeb1708751fef8039f3689807c018c Mon Sep 17 00:00:00 2001 From: MythEclipse Date: Sat, 13 Jun 2026 14:38:06 +0700 Subject: [PATCH] fix(gateway): voice speaking handler, metadata cache, transmitter backpressure MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - speakingHandler: check subscriptions.has() BEFORE voiceActiveUser broadcast - segment.ts: add LRU metadata cache (200 entries, avoids Discord API in hotpath) - transmitter.ts: add backpressure handling with drain queue/flush (PassThrough.write() return value was ignored — potential OOM under load) - cleanup backpressure queue on transmitter stop Co-Authored-By: Claude --- .../voice-recording/recorder/segment.ts | 15 +++++++++---- .../recorder/speakingHandler.ts | 7 ++++--- .../modules/voice-recording/transmitter.ts | 21 +++++++++++++++++-- 3 files changed, 34 insertions(+), 9 deletions(-) diff --git a/services/discord-gateway/src/modules/voice-recording/recorder/segment.ts b/services/discord-gateway/src/modules/voice-recording/recorder/segment.ts index 3439516..b89dc40 100644 --- a/services/discord-gateway/src/modules/voice-recording/recorder/segment.ts +++ b/services/discord-gateway/src/modules/voice-recording/recorder/segment.ts @@ -13,11 +13,15 @@ import type { RecordingSession } from "./sessionRecording.js"; import { uploadRecordingSegment } from "./uploader.js"; // --------------------------------------------------------------------------- -// Logger +// Logger & metadata cache // --------------------------------------------------------------------------- const logger = createChildLogger("voice-segment"); +/** LRU-ish cache: userId -> UserMetadata. Avoids Discord API calls in hotpath. */ +const metadataCache = new Map(); +const METADATA_CACHE_MAX = 200; + // --------------------------------------------------------------------------- // collectUserMetadata (was metadata.ts) // --------------------------------------------------------------------------- @@ -27,7 +31,8 @@ export async function collectUserMetadata( userId: string, channel: VoiceChannel, ): Promise { - logger.debug({ userId }, "Collecting user metadata"); + const cached = metadataCache.get(userId); + if (cached) return cached; const user = client.users.cache.get(userId) || @@ -52,7 +57,7 @@ export async function collectUserMetadata( position: role.position, })) ?? []; - return { + const result: UserMetadata = { userId, username, tag: user?.tag ?? "Unknown#0000", @@ -76,7 +81,9 @@ export async function collectUserMetadata( highestRole: roles[0] ?? null, joinedTimestamp: member?.joinedTimestamp ?? null, }; -} + + cacheMetadata(userId, result); + return result; // --------------------------------------------------------------------------- // Path helpers (was segment.ts) diff --git a/services/discord-gateway/src/modules/voice-recording/recorder/speakingHandler.ts b/services/discord-gateway/src/modules/voice-recording/recorder/speakingHandler.ts index 409c0d9..9517836 100644 --- a/services/discord-gateway/src/modules/voice-recording/recorder/speakingHandler.ts +++ b/services/discord-gateway/src/modules/voice-recording/recorder/speakingHandler.ts @@ -54,6 +54,10 @@ export function createSpeakingHandler( "Voice activity detected", ); + // Skip if user already has an active stream subscription + // (check BEFORE broadcast to avoid false positive events) + if (receiver.subscriptions.has(userId)) return; + // Notify webserver / WebSocket clients eventBroadcaster?.voiceActiveUser(userId, { username: userMetadata.username, @@ -61,9 +65,6 @@ export function createSpeakingHandler( speaking: true, }); - // Skip if user already has an active stream subscription - if (receiver.subscriptions.has(userId)) return; - // Ensure per-user recording directory const userDir = path.join(recordingsDir, userId); await fsPromises.mkdir(userDir, { recursive: true }).catch(() => { diff --git a/services/discord-gateway/src/modules/voice-recording/transmitter.ts b/services/discord-gateway/src/modules/voice-recording/transmitter.ts index b09e732..47daac3 100644 --- a/services/discord-gateway/src/modules/voice-recording/transmitter.ts +++ b/services/discord-gateway/src/modules/voice-recording/transmitter.ts @@ -21,6 +21,9 @@ export class VoiceTransmitter { private ffmpegProcess: ReturnType | null = null; private isActive = false; private readonly TRANSMIT_CHANNEL = BACKEND_VOICE_TRANSMIT; + /** Queue for PCM chunks when backpressure is active */ + private backpressureQueue: Buffer[] = []; + private draining = false; /** * Start listening for PCM audio data from Redis and stream to Discord @@ -130,8 +133,19 @@ export class VoiceTransmitter { const data = JSON.parse(message); if (data.type === "pcm" && data.buffer) { const pcmBuffer = Buffer.from(data.buffer, "base64"); - logger.debug({ bytes: pcmBuffer.length }, "Received PCM chunk"); - this.pcmStream.write(pcmBuffer); + const canContinue = this.pcmStream.write(pcmBuffer); + // Backpressure: queue until drain + if (!canContinue) { + this.draining = true; + this.pcmStream.once("drain", () => { + this.draining = false; + // Flush queued chunks + while (this.backpressureQueue.length > 0) { + const queued = this.backpressureQueue.shift()!; + if (!this.pcmStream.write(queued)) break; + } + }); + } } } catch (err) { logger.error({ error: err }, "Failed to process PCM data"); @@ -150,6 +164,9 @@ export class VoiceTransmitter { this.isActive = false; + this.backpressureQueue = []; + this.draining = false; + if (this.pcmStream) { this.pcmStream.end(); this.pcmStream = null;