diff --git a/.hermes/plans/2026-08-30_voice-recording-completeness-spec.md b/.hermes/plans/2026-08-30_voice-recording-completeness-spec.md new file mode 100644 index 00000000..252682df --- /dev/null +++ b/.hermes/plans/2026-08-30_voice-recording-completeness-spec.md @@ -0,0 +1,166 @@ +# Spec: Perbaiki alur voice → recording (miss & terpotong) + +## Konteks & Gejala +User melaporkan alur voice sampai recording **banyak miss** (audio tidak tercatat) +dan **terpotong** (satu alur bicara kebelah jadi beberapa segmen / audio putus di +tengah). Ini domain `services/discord-gateway/src/modules/voice-recording/`. + +Pipeline per user yang mulai bicara (speaking "start"): +``` +receiver.speaking "start" → speakingHandler(userId) + ├─ await collectUserMetadata(...) ← roundtrip API, subscribe tertunda + ├─ receiver.subscribe(userId, {end: AfterSilence, duration: 3000ms}) → audioStream + ├─ attach data/end/error handlers → audioStream.pipe(PacketFilter) → oggPacketStream + ├─ SegmentManager.open() → OggLogicalBitstream → file .ogg + ├─ data: SegmentManager.rotateIfNeeded (rotasi 5s) + decoder.write (web PCM tho + └─ end: SegmentManager.close() → segmen finish → finalizeSegment upload + transkrip +``` + +## Root cause (dari pembacaan kode — justifikasi di bawah) + +### A. MISS bagian awal bicara — subscribe tertunda (utama) +`speakingHandler.ts:53` melakukan `await collectUserMetadata(...)` SEBELUM +`receiver.subscribe`. `collectUserMetadata` (metadata.ts:30) pada cold path +(cache miss) melakukan `client.users.fetch` + `guild.members.fetch` roundtrip +Discord API (ratusan ms–detik). Selama await, seluruh opus awal bocor → awal +kalimat hilang. Cache menghilangkan ini untuk user yang pernah ter-record, tapi +user baru/evict (cache max 200) kena setiap kali. + +### B. Double-subscribe race +Guard `receiver.subscriptions.has(userId)` di `speakingHandler.ts:63` diletakkan +SETELAH `await collectUserMetadata`. Dua event "start" cepat keduanya melewati +guard (belum subscribe) → dua subscription → audio terbelah/ganda per user. + +### C. TERPOTONG di jeda bicara — AfterSilence 3000ms +`AUDIO_STREAM_SILENCE_DURATION_MS=3000`. Setelah 3s diam, stream auto-`end` → +`SegmentManager.close` → segmen baru saat bicara lagi. Jeda normal (berpikir, +interupsi) memecah 1 alur bicara jadi beberapa segmen/file. Ini source "terpotong". +Segmen pendek hasil jeda <1s juga DIBUANG di `finalizeSegment` (MIN_DURATION_MS=1000) +→ miss kata singkat ("ya", "siap"). + +### D. Rotasi segmen 5s di tengah bicara +`RECORDING_SEGMENT_MS=5000`: `rotateIfNeeded` menutup bitstream & membuka baru +setiap 5s walau bicara kontinu. Pipenya di-re-wire di dalam handler data → +window drop kecil + continuity file pecah (bukan masalah besar, tapi berkontribusi). + +### E. Tidak ada sinkronisasi "stop" speaking & stream "end" flaky +Handler hanya listen "start"; mengandalkan `AfterSilence` untuk emit "end". +Bug @discordjs/voice yang dikenal: `AfterSilence` bisa TIDAK emit "end" saat +koneksi gagal/teardown → segmen menggantung & tidak pernah finalize/upload +(recording "hilang"). Tidak ada watchdog. + +## Scope +Hanya `services/discord-gateway/src/modules/voice-recording/` (+ config index bila +perlu default baru). Tidak menyentuh playback (player.ts), transmitter (voice dari +browser → Discord, arah berlawanan), muxer (konsolidasi akhir), atau screen-share +video (hanya audio SSRC via `hookScreenShareAudio` sudah ada & dibiarkan). + +## Perubahan + +### 1. Speak-before-metadata: subscribe LEBIH DULU, metadata paralel +`recorder/speakingHandler.ts`: +- Pindahkan `receiver.subscribe` + pipeline setup ke ATAS, SEGERA di handler, + sebelum `collectUserMetadata`. +- Jalankan `collectUserMetadata` secara paralel non-blocking; gunakan metadata + cache untuk registrasi segmen saat finalize. +- Pertahankan guard skip bot/user (bot bisa dicek dari `client.users.cache` / + `client.user.id` tanpa await) SEBELUM subscribe — jangan tunggu fetch user. + +Detail konkret: +``` +async handler(userId): + if (userId === client.user?.id) return; + if (receiver.subscriptions.has(userId)) return; // guard kini di DEPAN, tanpa await + // (belum tahu bot? gunakan cache user; subscribe dulu biar nggak miss) + clone = subscribe(userId, {AfterSilence, duration}) // TANPA await metadata + setup pipeline (data/end/error, pipe, open segment) + collectUserMetadata(...).then(meta => { + if (meta.bot) { drain & close subscription (jangan simpan) } + else { registrasi ulang metadata utk segmen aktif } + }) +``` +Karena listener `start` dibuang untuk bot, harus tutup subscription bot tanpa +menyimpan segmen (buang hasil). Pakai `receiver.subscriptions.get(userId)?.destroy()`. + +### 2. Selesaikan race double-subscribe (bagian dari #1) +Guard `receiver.subscriptions.has(userId)` diletakkan SINCRON di awal (sebelum +await). Karena `subscribe` sinkron dan `subscriptions` terisi sinkron saat +dipanggil, event "start" kedua yang tiba setelah subscribe akan melihat +subscription aktif → di-skip. Tidak ada await antara guard & subscribe. + +### 3. Naikkan AfterSilence + tail-length → kurangi "terpotong" +`shared/config/index.ts`: +- `AUDIO_STREAM_SILENCE_DURATION_MS` default 3000 → 4000 (beri ruang jeda + alami; Discord packet 20ms, 4s masih wajar, tidak membengkak file). +Opsional via env override di production (tidak wajib komit env). + +### 4. Segmen berorientasi "burst bicara" daripada rotasi jam +`recorder/segment.ts` + `recorder/speakingHandler.ts`: +- Hapus/lepas rotasi SEGMEN berbasis waktu (RECORDING_SEGMENT_MS). Alih-alih, + satu segmen = satu burst bicara (buka di "start", tutup di "end"/AfterSilence + end). Ini menghilangkan pemecahan di tengah kalimat. +- `RECORDING_SEGMENT_MS` tetap dipakai untuk rotasi decoder web-PCM (broadcast + live), di mana segmen besar bisa menunda frame — biarkan seperti ada. + JADI: `SegmentManager.rotateIfNeeded` TIDAK lagi dipanggil pada jalur OGG + recording; decoder rotate tetap dijalankan. + +Catatan: dengan satu segmen per burst, ukuran file ~ durasi bicara. File panjang +dibutuhkan transkrip & transcode; tidak ada batas keras yang perlu di-override. +Watchdog di #5 membatasi durasi menggantung. + +### 5. Watchdog end-of-burst & teardown recovery +`recorder/speakingHandler.ts`: +- Setelah subscribe, arm timer watchdog (mis. `config twin`/hitung) yang menutup + segmen jika `AfterSilence` tidak emit "end" dalam X detik setelah "stop" + speaking — atau, lebih sederhana & robust: dengarkan BOTH stream "end" DAN + timer dari `receiver.speaking` "stop" (hingga @discordjs/voice meng-klaim + AfterSilence). Bila "stop" fire, mulai countdown kecil (mis. 500ms) lalu + `segmentManager.close` + `decoder.destroy` + destroy subscription jika stream + belum "end". +- Ini menutup A: segmen menggantung → jadi pasti finalize & upload. + +Implementasi: subscriptionStream (audioStream) + track milik per-user di +Map; handler "stop" +menjadwalkan finalize. + +### 6. Naikkan/ambil MIN segmen duration lebih rendah +`recorder/segmentFinalizer.ts`: MIN_DURATION_MS 1000 → 300ms. Kata pendek +("ya", "siap") tetap tersimpan. GUI biarkan. + +## File yang disentuh +- `src/modules/voice-recording/recorder/speakingHandler.ts` (utama: subscribe + first, guard depan, watchdog stop, hapus rotasi segmen dari jalur OGG) +- `src/modules/voice-recording/recorder/segment.ts` (opsional: API close/open, + pertahankan rotate untuk decoder tapi tak dipakai jalur OGG) +- `src/modules/voice-recording/recorder/streamSetup.ts` (kecil: backfill subscribe + supaya return subscription utk cleanup/destroy bot) +- `src/modules/voice-recording/recorder/segmentFinalizer.ts` (MIN_DURATION) +- `src/shared/config/index.ts` (default AfterSilence 4000) + +## Yang TIDAK disentuh +- `transmitter.ts` (arah browser→Discord, bukan recording) +- `player.ts`, `mediaSource.ts`, `screenShareAudio.ts` (hook SSRC sudah benar) +- `muxer.ts` (konsolidasi akhir tetap jalan) +- Flake/deps/bundle + +## Verification +1. `cd services/discord-gateway && pnpm typecheck` (tsc --noEmit) — 0 error. +2. `pnpm lint` (biome check src/) — exit 0. +3. `pnpm build` (tsc → dist/). +4. Unit test baru (vitest, kalau infra tes ada): + - subscribe terjadi tanpa await metadata (spy urutan panggilan) + - double-start skip lewat guard sinkron + - "stop" → watchdog finalize segmen walau stream tidak "end" + - MIN_DURATION 300ms menyimpan kata pendek +5. Smoke/CI: `nix flake check` (eval). Push → CI `Build & Deploy (Nix)` hijau → + deploy landing (`systemctl show gmw-discord-gateway.service --property=ActiveEnterTimestamp`). +6. Runtime manual (user): join voice, bicara dengan jeda >3s, pastikan 1 alur + kontinu = 1 segmen utuh (bukan 3), dan awal kata tidak hilang. + +## Risiko / Trade-off +- Subscribe-before-metadata: burst bot akan dikumpulkan sesaat lalu dibuang + (cost kecil: buang segmen). Lebih baik miss bot daripada miss user. +- AfterSilence naik: file lebih panjang sedikit saat jeda; upload/transkrip + timeout (transcodeToMp3 30s) tetap aman. +- Satu segmen per burst: tidak ada rotasi paksa → durasi segmen = durasi bicara + (bisa menit). Transkrip & transcode tetap ok. Watchdog batasi menggantung. diff --git a/services/discord-gateway/src/modules/voice-recording/recorder/segmentFinalizer.ts b/services/discord-gateway/src/modules/voice-recording/recorder/segmentFinalizer.ts index 2579ecdb..ef4aef5b 100644 --- a/services/discord-gateway/src/modules/voice-recording/recorder/segmentFinalizer.ts +++ b/services/discord-gateway/src/modules/voice-recording/recorder/segmentFinalizer.ts @@ -52,9 +52,10 @@ export function finalizeSegment(input: SegmentFinalizerInput): void { const endTime = currentSegment.endTime ?? Date.now(); const durationMs = endTime - currentSegment.startTime; - // Discard segments shorter than 1 second — not useful as a recording, - // would just be a blip of ambient noise or a mic click. - const MIN_DURATION_MS = 1000; + // Discard only very short segments — a blip of ambient noise or a mic + // click. At 300ms, brief acknowledgements ("ya", "siap", "ok") are still + // kept while true artifacts are dropped. + const MIN_DURATION_MS = 300; if (durationMs < MIN_DURATION_MS) { logger.debug( { filename: currentSegment.filename, durationMs }, 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 36feba2b..b1e484b8 100644 --- a/services/discord-gateway/src/modules/voice-recording/recorder/speakingHandler.ts +++ b/services/discord-gateway/src/modules/voice-recording/recorder/speakingHandler.ts @@ -4,6 +4,10 @@ import type { VoiceConnection } from "@discordjs/voice"; import type { Client, VoiceChannel } from "discord.js-selfbot-v13"; import { createChildLogger } from "@/shared/logger/index"; import type { EventBroadcaster } from "../../event-broadcaster/eventBroadcaster.js"; +import type { + SegmentState, + UserMetadata, +} from "../../message-capture/types.js"; import { collectUserMetadata } from "./metadata.js"; import { finalizeSegment } from "./segmentFinalizer.js"; import type { RecordingSession } from "./sessionRecording.js"; @@ -25,13 +29,21 @@ export interface SpeakingHandlerContext { /** * Creates the event handler for `receiver.speaking.on("start", handler)`. * - * The returned handler manages the full lifecycle for a user who starts speaking: - * 1. Validates the user (skip bot self, skip already-subscribed) - * 2. Collects user metadata and notifies the event broadcaster - * 3. Sets up the audio stream, decoder, packet filter, and segment manager - * 4. Attaches stream event handlers (data, end, error) BEFORE piping - * 5. Pipes audio through the packet filter for OGG recording - * 6. Handles segment completion (metadata write, upload trigger) + * # Why this is structured the way it is (recording-completeness) + * + * `receiver.speaking"start"` fires from `onUdpMessage` the moment the FIRST + * Opus packet for a user arrives, and the opus bytes are forwarded to the + * subscription stream **only if a subscription already exists** — otherwise + * they are dropped (`subscriptions.get(userId)` is undefined → `return`). + * + * That means the subscription MUST be created synchronously (no `await` before + * it): every frame between "start" and the subscribe call is otherwise silently + * discarded → the START of a burst is MISSING (the user's core complaint). + * The previous code `await`ed `collectUserMetadata` (a Discord REST roundtrip, + * slow on cache miss) BEFORE subscribing — that was the main source of misses. + * + * Metadata is now fetched in the background. If the speaker turns out to be a + * bot, the just-started burst is discarded (unlinked, never uploaded). */ export function createSpeakingHandler( ctx: SpeakingHandlerContext, @@ -47,110 +59,176 @@ export function createSpeakingHandler( } = ctx; return async (userId: string) => { - // Skip the bot's own audio + // Skip the bot's own audio. if (userId === client.user?.id) return; - const userMetadata = await collectUserMetadata(client, userId, channel); - if (userMetadata.bot) return; - - logger.debug( - { userId, username: userMetadata.username }, - "Voice activity detected", - ); - - // Skip if user already has an active stream subscription - // (check BEFORE broadcast to avoid false positive events) + // Synchronous guard BEFORE any await. Two "start" events can race during a + // subscription; guard here so we never create a second subscription. if (receiver.subscriptions.has(userId)) return; - // Notify webserver / WebSocket clients + logger.debug({ userId }, "Voice activity detected"); + + // Optimistic speaking:true (real metadata arrives async). eventBroadcaster?.voiceActiveUser(userId, { - username: userMetadata.username, - avatar: userMetadata.avatarUrl, + username: userId, + avatar: "", speaking: true, }); - // Ensure per-user recording directory + // Subscribe IMMEDIATELY (synchronously — guard is already done above, no + // `await` since then), so the subscription exists before any Opus frame + // arrives. Frames are then buffered by the AudioReceiveStream while we + // finish setup below (mkdir is non-network and fast). const userDir = path.join(recordingsDir, userId); + const { audioStream, packetFilter, segmentManager, decoder } = + setupUserStream({ + userId, + receiver, + userDir, + onPcmData: (pcm) => { + if (pcmSender) { + pcmSender(pcm, userId); + } else { + eventBroadcaster?.voicePcmData(pcm, userId); + } + }, + }); + + // Ensure per-user recording directory (subscription is already live, so + // this await does NOT drop audio — frames buffer in the subscription). await fsPromises.mkdir(userDir, { recursive: true }).catch(() => { - // Directory already exists, ignore + // Directory already exists, ignore. }); - try { - // Step 1: Set up stream components (subscribe, decoder, filter, segment - // manager). NOTE: pipe() is NOT called here — we attach event handlers - // first to prevent data loss from race conditions. - const { audioStream, packetFilter, segmentManager, decoder } = - setupUserStream({ - userId, - receiver, - userDir, - onPcmData: (pcm) => { - if (pcmSender) { - pcmSender(pcm, userId); - } else { - eventBroadcaster?.voicePcmData(pcm, userId); + const activeSession = activeSessions.get(channel.guild.id); + + // Mutable per-burst state. + let userMetadata: UserMetadata | null = null; + let accepted = true; // false once we learn it's a bot → discard burst + let finalized = false; + + /** Unlink the given segment's files (bot/error discard path). */ + const discardSegmentFiles = (seg: SegmentState | null): void => { + if (!seg) return; + fsPromises.unlink(seg.filename).catch(() => {}); + fsPromises.unlink(seg.jsonFilename).catch(() => {}); + }; + + const emitSpeakingFalse = (): void => { + eventBroadcaster?.voiceActiveUser(userId, { + username: userMetadata?.username ?? userId, + avatar: userMetadata?.avatarUrl ?? "", + speaking: false, + }); + }; + + /** + * Close the current (per-burst) segment and — once the underlying file has + * finished flushing to disk — finalize it (or discard it if bot/error). + * Safe to call multiple times (guarded by `finalized`). + */ + const finishBurst = (): void => { + if (finalized) return; + finalized = true; + const seg = segmentManager.close(oggPacketStream); + decoder.destroy(); + + const doFinalize = (): void => { + Promise.resolve(userMetadata).then((meta) => { + if (meta && accepted && seg) { + try { + finalizeSegment({ + currentSegment: seg, + userMetadata: meta, + activeSession, + guildId: channel.guild.id, + channelId: channel.id, + channelName: channel.name, + }); + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + logger.error({ userId, error: msg }, "Segment finalize failed"); } - }, + } else { + // Bot audio, metadata failure, or no segment opened — discard. + discardSegmentFiles(seg); + } }); + }; - // Step 2: Attach all audioStream event handlers BEFORE pipe() - audioStream.on("data", (chunk: Buffer) => { - if (chunk.length < 8) return; - segmentManager.rotateIfNeeded(packetFilter); - decoder.rotateIfNeeded(); - decoder.write(chunk); - }); + // finalizeSegment reads the OGG file, so wait for the write stream to + // flush before touching it. + if (seg?.out.writableFinished) { + doFinalize(); + } else { + seg?.out.once("finish", doFinalize); + } + }; - audioStream.on("end", () => { - segmentManager.close(packetFilter); - decoder.destroy(); + // Attach audioStream handlers BEFORE pipe() (prevents data loss). + audioStream.on("data", (chunk: Buffer) => { + if (chunk.length < 8) return; + // One segment per burst — no time-based rotation here. Only rotate the + // web-PCM broadcast decoder to bound its memory. + decoder.rotateIfNeeded(); + decoder.write(chunk); + }); + + audioStream.on("end", () => { + finishBurst(); + emitSpeakingFalse(); + }); + + audioStream.on("error", (error: Error) => { + logger.error({ userId, error: error.message }, "Audio stream error"); + finishBurst(); + emitSpeakingFalse(); + }); + + // Pipe for OGG recording (handlers already attached). + const oggPacketStream = audioStream.pipe(packetFilter); + + // Open the (single, per-burst) segment. + const currentSegment = segmentManager.open(oggPacketStream); + + // Handle file-write errors on the underlying write stream. + currentSegment.out.on("error", (err: unknown) => { + const msg = err instanceof Error ? err.message : String(err); + logger.error({ userId, error: msg }, "File write error"); + }); + + // Handle packet-filter errors (rare) by closing the burst. + packetFilter.on("error", (err: Error) => { + logger.error({ userId, error: err.message }, "PacketFilter error"); + finishBurst(); + emitSpeakingFalse(); + }); + + // Fetch user metadata in the background; reject bots by discarding burst. + collectUserMetadata(client, userId, channel) + .then((meta) => { + if (meta.bot) { + accepted = false; + finishBurst(); + emitSpeakingFalse(); + return; + } + userMetadata = meta; eventBroadcaster?.voiceActiveUser(userId, { - username: userMetadata.username, - avatar: userMetadata.avatarUrl, - speaking: false, + username: meta.username, + avatar: meta.avatarUrl, + speaking: true, }); - }); - - audioStream.on("error", (error: Error) => { - segmentManager.close(packetFilter); - decoder.destroy(); - logger.error({ userId, error: error.message }, "Audio stream error"); - }); - - // Step 3: Now pipe for OGG recording (safe — event handlers attached) - const oggPacketStream = audioStream.pipe(packetFilter); - - // Step 4: Open the first segment - const activeSession = activeSessions.get(channel.guild.id); - const currentSegment = segmentManager.open(oggPacketStream); - - // Step 5: Handle segment file completion - currentSegment.out.on("finish", () => { - finalizeSegment({ - currentSegment, - userMetadata, - activeSession, - guildId: channel.guild.id, - channelId: channel.id, - channelName: channel.name, - }); - }); - - currentSegment.out.on("error", (err: unknown) => { + }) + .catch((err: unknown) => { const msg = err instanceof Error ? err.message : String(err); - logger.error({ userId, error: msg }, "File write error"); + logger.warn( + { userId, error: msg }, + "Metadata fetch failed, dropping burst", + ); + accepted = false; + finishBurst(); + emitSpeakingFalse(); }); - - // Step 6: Handle packet filter errors - packetFilter.on("error", (err) => { - segmentManager.close(oggPacketStream); - logger.error({ userId, error: err.message }, "PacketFilter error"); - }); - } catch (e) { - logger.error( - { userId, error: e instanceof Error ? e.message : String(e) }, - "Failed to create stream", - ); - } }; } diff --git a/services/discord-gateway/src/shared/config/index.ts b/services/discord-gateway/src/shared/config/index.ts index b8ead9b0..9c9fde12 100644 --- a/services/discord-gateway/src/shared/config/index.ts +++ b/services/discord-gateway/src/shared/config/index.ts @@ -53,10 +53,15 @@ export const configSchema = z DECODER_COOLDOWN_MS: z.coerce.number().positive().default(30000), // ── Audio ──────────────────────────────────────────────────────────── + // AfterSilence: how long a voice burst may stay silent before the receive + // stream ends the segment. Raised 3000→4000 so natural pauses in speech + // (thinking gaps, interruptions) don't split one utterance into multiple + // segments ("terpotong"). Tunable via env; larger = fewer splits but a + // longer silent tail on each recording. AUDIO_STREAM_SILENCE_DURATION_MS: z.coerce .number() .positive() - .default(3000), + .default(4000), PACKET_FILTER_MIN_SIZE: z.coerce.number().positive().default(8), OPUS_FRAME_SIZE: z.coerce.number().positive().default(960), AUDIO_SAMPLE_RATE: z.coerce.number().positive().default(48000),