diff --git a/.hermes/plans/2026-08-30_video-record-receive-spec.md b/.hermes/plans/2026-08-30_video-record-receive-spec.md new file mode 100644 index 00000000..cde5f390 --- /dev/null +++ b/.hermes/plans/2026-08-30_video-record-receive-spec.md @@ -0,0 +1,66 @@ +# Spec: Record video (kamera/screenshare) orang lain — WebRTC receive (Request 2) + +Date: 2026-08-30. Status: IN PROGRESS (Phase A). + +## Why this is hard (ground truth, verified from @discordjs/voice 0.19.2 source) +`VoiceReceiver.onUdpMessage` (dist/index.mjs:2059) drops EVERY non-opus RTP +packet at line 2068: `if ((msg[1] & 127) !== RTP_OPUS_PAYLOAD_TYPE) return;`. +So video (kamera H264? actually Discord uses VP8/H264; screenshare combines with +video SSRC) is decrypted-capable but never forwarded. `receiver.parsePacket` +(2033) DOES decrypt any payload type generically (audio + video) using +`connectionData.{encryptionMode, nonceBuffer, secretKey}` — the only audio gate +is the opus check inside onUdpMessage. + +=> FIX: wrap `receiver.onUdpMessage` (like screenShareAudio.ts already does for +screen-share AUDIO SSRCs): for packets whose payload type is a VIDEO type +(payload 96 VP8, 101/102 H264, 106/116/126/127 AV1, VP9 98...), call +`receiver.parsePacket(...)` myself to decrypt, then depacketize + write frames. +Delegate opus (120) to the original handler. Delegate audio to original. + +## Audio already works (screenShareAudio.ts). We add VIDEO. + +## Science-of-the-changes below. + +## Phase A — capture + decrypt + depacketize to AnnexB h264 (THIS change) +Files (new): `src/modules/voice-recording/videoReceiver.ts` +- Hook into `recorder.startRecording` alongside `hookScreenShareAudio`. +- Wrap `receiver.onUdpMessage`: + - read ssrc = msg.readUInt32BE(8); userData = receiver.ssrcMap.get(ssrc) + - if payload type is video AND we have a "watching" subscription for that user + (videoSSRC present), decrypt via receiver.parsePacket(...), then: + - H264 (101/102 + payload 120 not): strip RTP header, reassemble FU-A + fragments into AnnexB NALs (start-code prefixed), buffer until we have + a full access unit (keyframe SPS/PPS/IDR or slices), append to a per- + user-per-burst `.h264` file. + - else delegate to original onUdpMessage. +- Watch `receiver.ssrcMap` "create"/"update" for `videoSSRC !== undefined` → + signal a video burst started for that user (like screenShareAudio does). +- Per-user video files written to `config.RECORDINGS_DIR//video-.h264`. +- Guard: skip bot's own video (client.user.id). + +Dependencies: NO new npm deps for Phase A (only crypto already in +@discordjs/voice via parsePacket + Buffer). ffmpeg-headless (already in Nix +buildInputs) used in Phase B for decode+mux. + +## Phase B — decode + mux to playable MP4/WebM (follow-up) +- Pipe AnnexB h264 → `ffmpeg -f h264 -i pipe:0 -c copy out.mp4` (or transcode). +- Reuse muxer patterns from voice (muxer.ts) + session recording. +- Persist in voice_recordings / new video_recordings table; upload via teleuploader. + +## Phase C — frontend playback + session grouping (follow-up) +- Backend oRPC list video files; FE video player, group by call session like audio. + +## Verification +- Phase A: join voice, have a member screen-share/camera, confirm `.h264` file + grows with NAL frames + keyframes; journal shows "video burst" logs. +- Run vitest unit: RTP header strip + FU-A reassembly gives correct bytes. + +## Open questions / risks +- Discord codec for camera = H264(101/102); screenshare uses H264 (101/103?) + and can also be VP8/VP9. Handle H264 first (depacketize proven), VP8/VP9 in + Phase B via ffmpeg RTP input. +- Encryption: DAVE (dave_protocol_version) adds a session layer; parsePacket + already applies daveSession.decrypt for audio — we must call the SAME + parsePacket path so DAVE/encryption is handled identically. +- ssrc↔user mapping during a broadcast: videoSSRC is in ssrcMap after the + voice state; may need the STREAM_CREATE network events to key reliably. diff --git a/services/discord-gateway/src/modules/voice-recording/recorder.ts b/services/discord-gateway/src/modules/voice-recording/recorder.ts index f82c82f6..95c64322 100644 --- a/services/discord-gateway/src/modules/voice-recording/recorder.ts +++ b/services/discord-gateway/src/modules/voice-recording/recorder.ts @@ -19,6 +19,7 @@ import { } from "./recorder/sessionRecording.js"; import { createSpeakingHandler } from "./recorder/speakingHandler.js"; import { hookScreenShareAudio } from "./screenShareAudio.js"; +import { hookVideoReceiver } from "./videoReceiver.js"; const logger = createChildLogger("recorder"); @@ -192,6 +193,16 @@ export async function startRecording( // to discover and register those SSRCs. hookScreenShareAudio(receiver, speakingHandler); + // ── Video capture (camera + screen share) ─────────────────────────── + // @discordjs/voice only decrypts/forwards audio (opus). Video RTP (H264, + // VP8/VP9/AV1) arrives on the same UDP socket but is dropped by the opus + // gate. hookVideoReceiver wraps onUdpMessage, decrypts video payload types + // with the connection secret key (reusing receiver.parsePacket), then + // depacketizes H264 to AnnexB and writes a raw .h264 stream per user. + // Order matters: video hook must be installed AFTER the screen-share audio + // hook so non-video packets delegate down the existing wrapper chain. + hookVideoReceiver(receiver, client, recordingsDir); + // Handle unexpected disconnection connection.on(VoiceConnectionStatus.Disconnected, async () => { if (config.VERBOSE) { diff --git a/services/discord-gateway/src/modules/voice-recording/videoReceiver.ts b/services/discord-gateway/src/modules/voice-recording/videoReceiver.ts new file mode 100644 index 00000000..d9d2bcbf --- /dev/null +++ b/services/discord-gateway/src/modules/voice-recording/videoReceiver.ts @@ -0,0 +1,369 @@ +import { createWriteStream, mkdirSync, type WriteStream } from "node:fs"; +import path from "node:path"; +import type { VoiceReceiver } from "@discordjs/voice"; +import type { Client } from "discord.js-selfbot-v13"; +import { createChildLogger } from "@/shared/logger/index"; + +const logger = createChildLogger("video-receiver"); + +// ─── Discord RTP video payload types (real Discord clients) ─────────────── +// H264 is what Discord uses for the user's CAMERA; screen-share (GoLive) also +// uses H264. VP8/VP9/AV1 appear for some share types / forced codecs. +export const VIDEO_PAYLOAD_TYPES = new Set([ + 96, 98, 101, 102, 106, 116, 126, 127, +]); + +// H264 NAL unit types (single byte, after the 1-byte NAL header) +const NAL_SPS = 7; +const NAL_PPS = 8; +const NAL_IDR = 5; + +// RTP H264 packet types (byte 0 of the payload) +const _H264_SINGLE_NAL = 1; // 0-23 → single NAL unit, type = nal_unit_type +const H264_STAP_A = 24; +const H264_FU_A = 28; + +/** + * Pure RTP-H264 depacketizer → AnnexB (prefix NALs with 00 00 00 01 start codes). + * + * Input: RTP payload bytes (after the 12-byte RTP header, after decryption). + * Output: a list of AnnexB NAL buffers (each with start code) ready to append + * to a `.h264` elementary stream file. + * + * Reassembles FU-A fragmented NALs and STAP-A aggregated NALs into standalone + * NALs with start codes. Fragments that don't form a complete NAL are buffered. + */ +export class H264Depacketizer { + /** True after we've seen an SPS/PPS/IDR so we only start writing on a keyframe. */ + private sawConfig = false; + /** Buffer for an in-progress FU-A NAL. */ + private fuBuffer: number[] = []; + + /** Handle one RTP payload; returns AnnexB NAL buffers to write (may be []). */ + push(payload: Buffer): Buffer[] { + const out: Buffer[] = []; + if (payload.length === 0) return out; + + const packetType = payload[0] & 0x1f; + + if (packetType >= 1 && packetType <= 23) { + // Single NAL unit. + const nalType = payload[0] & 0x1f; + this.maybeEnableConfig(nalType); + const annexB = this.toAnnexB(payload); + if (this.sawConfig && annexB) out.push(annexB); + return out; + } + + if (packetType === H264_FU_A) { + // Fragmentation unit: payload[0] = FU indicator (F/NRI/type=28), + // payload[1] = FU header (S/E/R/type). + const fuHeader = payload[1]; + const start = (fuHeader & 0x80) !== 0; + const end = (fuHeader & 0x40) !== 0; + const nalType = fuHeader & 0x1f; + + if (start) { + // First fragment: reconstruct the full NAL header by preserving the + // F and NRI bits from the FU indicator (payload[0]) OR'ed with the + // NAL unit type from the FU header (payload[1] & 0x1F). + const nalHeader = (payload[0] & 0xe0) | nalType; + this.fuBuffer = [nalHeader, ...payload.subarray(2)]; + } else { + // Continuation. (Ignore if no buffer — a lost first fragment.) + if (this.fuBuffer.length === 0) { + // We missed the start; the NAL type unknown, skip this fragment. + return out; + } + if (payload.length > 2) { + this.fuBuffer.push(...payload.subarray(2)); + } + } + + if (end) { + const nal = Buffer.from(this.fuBuffer); + this.fuBuffer = []; + this.maybeEnableConfig(nalType); + if (this.sawConfig && nal.length > 0) { + out.push(this.toAnnexB(nal)); + } + } + return out; + } + + if (packetType === H264_STAP_A) { + // Aggregation packet: list of (2-byte NAL size, NAL data). + let offset = 1; + while (offset < payload.length) { + if (offset + 2 > payload.length) break; + const naluLength = payload.readUInt16BE(offset); + offset += 2; + if (offset + naluLength > payload.length) break; + const nalu = payload.subarray(offset, offset + naluLength); + offset += naluLength; + const nalType = nalu[0] & 0x1f; + this.maybeEnableConfig(nalType); + const annexB = this.toAnnexB(nalu); + if (this.sawConfig && annexB) out.push(annexB); + } + } + return out; + } + + private maybeEnableConfig(nalType: number): void { + if (nalType === NAL_SPS || nalType === NAL_PPS || nalType === NAL_IDR) { + this.sawConfig = true; + } + } + + private toAnnexB(nal: Buffer): Buffer { + const startCode = Buffer.from([0x00, 0x00, 0x00, 0x01]); + return Buffer.concat([startCode, nal]); + } + + reset(): void { + this.sawConfig = false; + this.fuBuffer = []; + } +} + +// ─── Per-user video burst state ──────────────────────────────────────────── + +interface VideoBurst { + userId: string; + depacketizer: H264Depacketizer; + /** video SSRC this burst is sourced from (may be camera or screenshare). */ + ssrc: number; + createdAt: number; + filePath: string; + out: WriteStream; + /** count of AnnexB NAL bytes written (for logging/health). */ + bytesWritten: number; + lastPacketAt: number; +} + +/** + * Listens for other users' VIDEO RTP (camera + screen share) on the same UDP + * socket @discordjs/voice receives audio from, decrypts each video packet with + * the connection's secret key (via receiver.parsePacket), depacketizes H264 to + * AnnexB, and writes a raw `.h264` elementary stream per user-per-burst. + * + * Wraps `receiver.onUdpMessage` like screenShareAudio.ts. Packets with a video + * payload type are consumed here; ALL other packets (opus mic + screen-share + * audio) are delegated to the current handler (the existing wrapper chain). + */ +export function hookVideoReceiver( + receiver: VoiceReceiver, + client: Client, + recordingsDir: string, +): void { + const original = receiver.onUdpMessage; + const bursts = new Map(); + + // The UDP socket delivers video on SSRCs that are NOT the audioSSRC key of + // ssrcMap. Build a videoSSRC→userId index from VoiceUserData as it updates so + // we can attribute packets without SSRC-proximity guessing. + const videoSsrcToUser = new Map(); + + receiver.ssrcMap.on("create", (data) => { + if (data.videoSSRC !== undefined) { + videoSsrcToUser.set(data.videoSSRC, data.userId); + logger.info( + { userId: data.userId, videoSSRC: data.videoSSRC }, + "Video SSRC appeared (camera/screenshare start)", + ); + } + }); + receiver.ssrcMap.on("update", (_old, neu) => { + if (neu.videoSSRC !== undefined) { + videoSsrcToUser.set(neu.videoSSRC, neu.userId); + } + }); + + /** Close a burst's file (flush + end). */ + const closeBurst = (userId: string): void => { + const burst = bursts.get(userId); + if (!burst) return; + bursts.delete(userId); + try { + burst.out.end(); + } catch { + // ignore write errors on close + } + logger.info( + { + userId, + ssrc: burst.ssrc, + bytes: burst.bytesWritten, + durationMs: Date.now() - burst.createdAt, + }, + "Video burst closed", + ); + }; + + const openBurst = (userId: string, ssrc: number): VideoBurst | null => { + try { + const userDir = path.join(recordingsDir, userId); + mkdirSync(userDir, { recursive: true }); + const filePath = path.join(userDir, `video-${ssrc}-${Date.now()}.h264`); + const out = createWriteStream(filePath); + const burst: VideoBurst = { + userId, + ssrc, + depacketizer: new H264Depacketizer(), + createdAt: Date.now(), + filePath, + out, + bytesWritten: 0, + lastPacketAt: Date.now(), + }; + bursts.set(userId, burst); + logger.info({ userId, ssrc, filePath }, "Video burst opened"); + return burst; + } catch (err) { + logger.warn( + { userId, ssrc, err: err instanceof Error ? err.message : String(err) }, + "Failed to open video burst file", + ); + return null; + } + }; + + receiver.onUdpMessage = (msg: Buffer) => { + if (msg.length <= 8) { + original.call(receiver, msg); + return; + } + + const payloadType = msg[1] & 127; + if (!VIDEO_PAYLOAD_TYPES.has(payloadType)) { + // Not video — let the existing handler (audio + screen-share audio) take it. + original.call(receiver, msg); + return; + } + + // ── It's a video RTP packet. Decrypt + depacketize + write. ── + const ssrc = msg.readUInt32BE(8); + + // Attribute to a user. + let userId = videoSsrcToUser.get(ssrc); + if (!userId) { + // Fallback: maybe we registered the SSRC under an audio entry via the + // screen-share-audio hook, or proximity. Skip bot's own stream. + const candidate = inferVideoOwner(msg, receiver); + if (candidate && candidate.userId !== client.user?.id) { + userId = candidate.userId; + videoSsrcToUser.set(ssrc, candidate.userId); + } + } + if (!userId || userId === client.user?.id) { + return; // unknown or self — ignore + } + + // Decrypt the payload using @discordjs/voice's parser (handles DAVE + + // encryption the same way it does for audio). + let decrypted: Buffer; + try { + const receiverAny = receiver as unknown as { + parsePacket: ( + buffer: Buffer, + mode: string, + nonce: Buffer, + secretKey: Uint8Array, + userId: string, + ) => Buffer; + connectionData: { + encryptionMode?: string; + nonceBuffer?: Buffer; + secretKey?: Uint8Array; + }; + }; + const { encryptionMode, nonceBuffer, secretKey } = + receiverAny.connectionData; + if (!encryptionMode || !nonceBuffer || !secretKey) return; + // decrypted payload = RTP header already stripped by parsePacket + decrypted = receiverAny.parsePacket( + msg, + encryptionMode, + nonceBuffer, + secretKey, + userId, + ); + } catch { + // Decryption failed — transient; skip this packet. + return; + } + + // Depacketize the (video) RTP payload → AnnexB. + let burst = bursts.get(userId); + if (!burst) { + const opened = openBurst(userId, ssrc); + if (!opened) return; + burst = opened; + } + burst.lastPacketAt = Date.now(); + + let nals: Buffer[]; + try { + nals = burst.depacketizer.push(decrypted); + } catch { + return; + } + if (nals.length === 0) return; + + for (const nal of nals) { + burst.bytesWritten += nal.length; + if (!burst.out.write(nal)) { + burst.out.once("drain", () => {}); + } + } + }; + + // If a user's video SSRC disappears, close their burst after a short grace so + // the tail flushes. + setInterval(() => { + const now = Date.now(); + for (const [userId, burst] of bursts) { + if (now - burst.lastPacketAt > 5000) { + closeBurst(userId); + } + } + }, 2000); + // Prevent the interval from keeping the process alive. + // (The gateway owns the process; a single 2s timer is negligible.) +} + +/** + * Infer which user owns a video SSRC by proximity to their known audioSSRC + * (Discord allocates a user's video SSRC near their audio SSRC). + */ +function inferVideoOwner( + msg: Buffer, + receiver: VoiceReceiver, +): { userId: string } | null { + const ssrc = msg.readUInt32BE(8); + try { + const map = getSsrcInternalMap(receiver.ssrcMap); + if (!map) return null; + for (const [, data] of map.entries()) { + if (data.audioSSRC && Math.abs(data.audioSSRC - ssrc) < 400_000) { + return { userId: data.userId }; + } + } + } catch { + // ignore + } + return null; +} + +/** Access the private `_map` inside SSRCMap (mirrors screenShareAudio.ts). */ +function getSsrcInternalMap( + ssrcMap: VoiceReceiver["ssrcMap"], +): + | Map + | undefined { + const asAny = ssrcMap as unknown as Record; + const m1 = asAny._map; + if (m1 instanceof Map) return m1 as Map; + return undefined; +} diff --git a/services/discord-gateway/tests/videoReceiver.test.ts b/services/discord-gateway/tests/videoReceiver.test.ts new file mode 100644 index 00000000..03e5ccbd --- /dev/null +++ b/services/discord-gateway/tests/videoReceiver.test.ts @@ -0,0 +1,62 @@ +import { describe, expect, it } from "vitest"; +import { H264Depacketizer } from "../src/modules/voice-recording/videoReceiver.js"; + +// Build a single-NAL RTP payload: [NAL header, ...data] +function singleNal(nalHeader: number, data: number[] = []): Buffer { + return Buffer.from([nalHeader, ...data]); +} + +describe("H264Depacketizer", () => { + it("prefixes a single NAL with an AnnexB start code once config seen", () => { + const d = new H264Depacketizer(); + // SPS first (0x67 = nal ref idc 3, type 7) + const sps = d.push(singleNal(0x67, [1, 2, 3])); + expect(sps.length).toBe(1); + expect(sps[0].subarray(0, 4)).toEqual(Buffer.from([0, 0, 0, 1])); + expect(sps[0][4]).toBe(0x67); // start code + original NAL + }); + + it("drops NALs before seeing SPS/PPS/IDR (waits for a keyframe)", () => { + const d = new H264Depacketizer(); + // A non-IDR slice (type 1) before any config → dropped + const res = d.push(singleNal(0x21, [9, 9, 9])); + expect(res.length).toBe(0); + }); + + it("emits frames after config", () => { + const d = new H264Depacketizer(); + d.push(singleNal(0x67, [])); // SPS + const slice = d.push(singleNal(0x65, [1, 2, 3, 4])); // IDR (type 5) + expect(slice.length).toBe(1); + expect(slice[0].equals(Buffer.from([0, 0, 0, 1, 0x65, 1, 2, 3, 4]))).toBe( + true, + ); + }); + + it("reassembles FU-A fragments into one NAL", () => { + const d = new H264Depacketizer(); + d.push(singleNal(0x67, [])); // SPS config + // FU-A: payload[0]=FU indicator (0x7C=type28), payload[1]=FU header + // start fragment: S=1 E=0, NAL type 5 => 0x85 ; data [10, 11] + const start = Buffer.from([0x7c, 0x85, 10, 11]); + const middle = Buffer.from([0x7c, 0x05, 12, 13]); + const end = Buffer.from([0x7c, 0x45, 14, 15]); // E=1 + + expect(d.push(start).length).toBe(0); // not complete yet + expect(d.push(middle).length).toBe(0); + const final = d.push(end); + expect(final.length).toBe(1); + // NAL header reconstructed as (0x7C & 0xE0) | 5 = 0x65, then data 10..15 + expect( + final[0].equals(Buffer.from([0, 0, 0, 1, 0x65, 10, 11, 12, 13, 14, 15])), + ).toBe(true); + }); + + it("first fragment of a burst with unknown start is dropped gracefully", () => { + const d = new H264Depacketizer(); + d.push(singleNal(0x67, [])); // SPS + // A continuation fragment with no buffer → ignored, no crash + const orphan = Buffer.from([0x7c, 0x05, 99]); // S=0 + expect(d.push(orphan).length).toBe(0); + }); +});