From 0ea76a8373827c7974198c04e1161e3d450b185a Mon Sep 17 00:00:00 2001 From: asepharyana Date: Mon, 31 Aug 2026 13:46:24 +0700 Subject: [PATCH] feat(gateway): DAVE-capable stream-watch video receive (Phase D) Replace the dead selfbot-v13 video path (WS 4017 DAVE). videoRecorder now delegates to a new streamWatchReceiver that: - sends STREAM_WATCH (op 20) on voiceState.streaming - opens a @discordjs/voice Networking to the watch RTC (STREAM_CREATE + STREAM_SERVER_UPDATE) with DAVE enabled - decrypts H264 via Davey MediaType.VIDEO, depacketizes + muxes to mp4 - tears down on streaming-stop / leave / untrack Remove ensureSelfbotVoice/createVideoStream/joinStreamConnection (dead). recorder.ts no longer fires the futile eager selfbot join. --- .../src/modules/voice-recording/recorder.ts | 13 +- .../voice-recording/streamWatchReceiver.ts | 387 ++++++++++++++++++ .../modules/voice-recording/videoRecorder.ts | 331 +++------------ .../tests/videoRecorder.test.ts | 163 +++----- 4 files changed, 489 insertions(+), 405 deletions(-) create mode 100644 services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts diff --git a/services/discord-gateway/src/modules/voice-recording/recorder.ts b/services/discord-gateway/src/modules/voice-recording/recorder.ts index 1b46e673..cd90434c 100644 --- a/services/discord-gateway/src/modules/voice-recording/recorder.ts +++ b/services/discord-gateway/src/modules/voice-recording/recorder.ts @@ -20,11 +20,7 @@ import { import { createSpeakingHandler } from "./recorder/speakingHandler.js"; import { hookScreenShareAudio } from "./screenShareAudio.js"; import { hookVideoReceiver } from "./videoReceiver.js"; -import { - ensureSelfbotVoice, - trackChannel, - untrackChannel, -} from "./videoRecorder.js"; +import { trackChannel, untrackChannel } from "./videoRecorder.js"; const logger = createChildLogger("recorder"); @@ -115,13 +111,6 @@ export async function startRecording( client: Client, channel: VoiceChannel, ): Promise { - // Establish the SELFbot voice connection FIRST (video receive). It MUST ride - // the bot's FRESH join so Discord emits VOICE_SERVER_UPDATE and the selfbot - // VoiceConnection authenticates. (Placing it after the @discordjs/voice join - // is Ready — a lazy re-join — times out and video capture fails.) Best-effort, - // fire-and-forget so it never blocks the audio join below. - void ensureSelfbotVoice(channel).catch(() => {}); - const connection = joinVoiceChannel({ channelId: channel.id, guildId: channel.guild.id, diff --git a/services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts b/services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts new file mode 100644 index 00000000..7d7939a5 --- /dev/null +++ b/services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts @@ -0,0 +1,387 @@ +/** + * DAVE-capable stream-watch video receiver for OTHER members (camera + screen + * share). Replaces the dead selfbot-v13 path (which Discord rejects with WS + * close 4017 "E2EE/DAVE protocol required" because that lib's voice stack + * predates DAVE). + * + * HOW IT WORKS (verified architecture, 2026-08-31): + * - Watching a member's stream is a SEPARATE RTC connection, not the guild + * audio socket. We send the gateway signal STREAM_WATCH (op 20) with + * stream_key `guild:::`; Discord replies with gateway events + * STREAM_CREATE (d.rtc_server_id) + STREAM_SERVER_UPDATE (d.token, + * d.endpoint) that together describe a dedicated voice RTC for the watch. + * - We open a `@discordjs/voice` `Networking` connection pointed at that RTC + * (endpoint/token/serverId = rtc_server_id + the bot's sessionId). `Networking` + * performs the FULL DAVE handshake (identify with max_dave_protocol_version, + * IP discovery, select-protocol, DAVE MLS transitions) — the identical path + * that already works for audio. It is a public export. + * - On Ready, `Networking` exposes `state.udp` (VoiceUDPSocket). We hook its + * `message` event. Video RTP (H264) arrives there; we decrypt via the DAVE + * session with `MediaType.VIDEO` (the wrapper's default is AUDIO), then feed + * the existing `H264Depacketizer` → AnnexB `.h264` → `muxToMp4`. + * + * All failures are best-effort and NEVER break the gateway's audio recording. + */ +import { + createWriteStream, + promises as fsPromises, + type WriteStream, +} from "node:fs"; +import path from "node:path"; +import { Networking, type NetworkingState } from "@discordjs/voice"; +import type Davey from "@snazzah/davey"; +import { MediaType } from "@snazzah/davey"; +import type { Client, VoiceChannel } from "discord.js-selfbot-v13"; +import { createChildLogger } from "@/shared/logger/index"; +import { H264Depacketizer, muxToMp4 } from "./videoReceiver.js"; + +const logger = createChildLogger("stream-watch"); + +/** A watched user's live DAVE stream connection + file state. */ +interface WatchState { + uid: string; + guildId: string; + channelId: string; + net: Networking; + depacketizer: H264Depacketizer; + filePath: string; + out: WriteStream; + bytesWritten: number; + lastPacketAt: number; + startedAt: number; +} + +const watches = new Map(); // key `${guildId}:${uid}` +const pendingServerId = new Map(); // key `${guildId}:${uid}` -> rtc_server_id (from STREAM_CREATE) + +let _client: Client | undefined; +let _rawAttached = false; +let _recordingsDir = "/var/lib/gmw/recordings"; + +export function setStreamWatchRecordingsDir(dir: string): void { + _recordingsDir = dir; +} + +/** Test-only reset. */ +export function __resetStreamWatchState(): void { + for (const [, w] of watches) closeWatch(w); + watches.clear(); + pendingServerId.clear(); + _rawAttached = false; + _client = undefined; +} + +/** Bind the selfbot client (needed for STREAM_WATCH + raw). Idempotent. */ +export function setStreamWatchClient(client: Client | undefined): void { + _client = client; + if (!client || _rawAttached) return; + _rawAttached = true; + client.on("raw", (packet: unknown) => { + try { + handleRaw(packet); + } catch (err) { + logger.warn( + { err: err instanceof Error ? err.message : String(err) }, + "stream-watch raw handler error (ignoring)", + ); + } + }); +} + +const streamKeyFor = (channelId: string, guildId: string, uid: string) => + `guild:${guildId}:${channelId}:${uid}`; + +/** Parse `guild:::` back into parts. */ +function parseStreamKey( + key: string, +): { guildId: string; channelId: string; uid: string } | null { + const m = /^guild:(\d+):(\d+):(\d+)$/.exec(key); + if (!m) return null; + return { guildId: m[1], channelId: m[2], uid: m[3] }; +} + +function handleRaw(packet: unknown): void { + const p = packet as { t?: string; d?: Record } | null; + if (!p?.t || !p?.d) return; + const t = p.t; + const streamKey = p.d.stream_key as string | undefined; + if (!streamKey) return; + const key = parseStreamKey(streamKey); + if (!key) return; + const watchKey = `${key.guildId}:${key.uid}`; + + if (t === "STREAM_CREATE") { + const rtc = String(p.d.rtc_server_id ?? ""); + pendingServerId.set(watchKey, rtc); + logger.info( + { guildId: key.guildId, userId: key.uid, rtcServerId: rtc }, + "STREAM_CREATE received (watch RTC authorized)", + ); + } else if (t === "STREAM_SERVER_UPDATE") { + const endpoint = String(p.d.endpoint ?? ""); + const token = String(p.d.token ?? ""); + const rtc = + pendingServerId.get(watchKey) ?? String(p.d.rtc_server_id ?? ""); + logger.info( + { guildId: key.guildId, userId: key.uid, endpoint, rtcServerId: rtc }, + "STREAM_SERVER_UPDATE received — connecting DAVE watch", + ); + connectWatch(key.guildId, key.channelId, key.uid, endpoint, token, rtc); + } else if (t === "STREAM_DELETE") { + logger.info( + { guildId: key.guildId, userId: key.uid }, + "STREAM_DELETE received — closing watch", + ); + pendingServerId.delete(watchKey); + closeWatchByKey(watchKey); + } else if (t === "STREAM_UPDATE") { + // paused / resumed — ignore + } +} + +/** + * Initiate watching `userId`'s stream in `channel`. Sends STREAM_WATCH and + * waits for STREAM_CREATE / STREAM_SERVER_UPDATE (handled by handleRaw) to + * actually open the DAVE connection. Best-effort. + */ +export function startStreamWatch(channel: VoiceChannel, userId: string): void { + if (!_client) { + logger.warn("startStreamWatch: no client set"); + return; + } + const watchKey = `${channel.guild.id}:${userId}`; + if (watches.has(watchKey)) return; // already watching + const sk = streamKeyFor(channel.id, channel.guild.id, userId); + logger.info( + { userId, guildId: channel.guild.id, channelId: channel.id, streamKey: sk }, + "Sending STREAM_WATCH to authorize video receive", + ); + try { + ( + _client as unknown as { ws: { broadcast: (d: unknown) => unknown } } + ).ws.broadcast({ + op: 20, // STREAM_WATCH + d: { stream_key: sk }, + }); + } catch (err) { + logger.warn( + { err: err instanceof Error ? err.message : String(err) }, + "Failed to send STREAM_WATCH", + ); + } +} + +function connectWatch( + guildId: string, + channelId: string, + uid: string, + endpoint: string, + token: string, + serverId: string, +): void { + const watchKey = `${guildId}:${uid}`; + if (watches.has(watchKey)) return; + const botId = _client?.user?.id; + if (!botId) return; + + // The bot's ACTIVE voice-session id (from the selfbot client's voice manager — + // the same session the guild @discordjs/voice connection rides). + const voiceManager = ( + _client as unknown as { + voice?: { connection?: { authentication?: { sessionId?: string } } }; + } + ).voice; + const sessionId = voiceManager?.connection?.authentication?.sessionId; + + logger.info( + { userId: uid, endpoint, serverId, hasSession: !!sessionId }, + "Opening DAVE Networking to watch RTC", + ); + + const net = new Networking( + { + endpoint, + serverId, + userId: botId, + sessionId: sessionId ?? "none", + token, + channelId, + }, + { daveEncryption: true }, + ); + + const uidProps = uid; + const netAny = net as unknown as { + state: NetworkingState & { + udp?: { + on: (e: string, cb: (m: Buffer) => void) => void; + off: (e: string, cb: (m: Buffer) => void) => void; + }; + dave?: { session?: Davey.DAVESession }; + }; + }; + + const onMessage = (msg: Buffer): void => { + handleUdpMessage(watchKey, netAny, uidProps, msg); + }; + + net.on("stateChange", (newState: unknown) => { + const s = newState as { code?: number }; + const code = s?.code; + if (code === 4 /* Ready */) { + const udp = netAny.state.udp; + if (udp) { + udp.on("message", onMessage); + logger.info( + { userId: uid }, + "DAVE watch READY — listening for video RTP", + ); + } else { + logger.warn({ userId: uid }, "watch ready but no UDP socket"); + } + } + }); + net.on("error", (err) => { + logger.warn( + { userId: uid, err: err instanceof Error ? err.message : String(err) }, + "watch Networking error", + ); + closeWatchByKey(watchKey); + }); + net.on("close", () => { + closeWatchByKey(watchKey); + }); + + watches.set(watchKey, { + uid, + guildId, + channelId, + net, + depacketizer: new H264Depacketizer(), + filePath: "", + out: undefined as unknown as WriteStream, + bytesWritten: 0, + lastPacketAt: Date.now(), + startedAt: Date.now(), + }); +} + +function handleUdpMessage( + watchKey: string, + net: { state: { dave?: { session?: Davey.DAVESession } } }, + uid: string, + msg: Buffer, +): void { + if (msg.length <= 12) return; + const payloadType = msg[1] & 127; + // Opus = 120. Anything else on the video/watch socket is video (H264 etc.). + // (We only watch a stream, so deliver everything non-opus to the video path.) + if (payloadType === 120) return; + const watch = watches.get(watchKey); + if (!watch) return; + // RTP header (12 bytes) + optional CSRC/extensions. For DAVE video we pass the + // RTP payload (after the header) to the session, matching audio's parsePacket. + // Minimal header parse: assume 12-byte header (Discord video RTP uses no header + // extension for this stream, matching the Discord-RE reference SDP). + const rtpPayload = msg.subarray(12); + const dave = net.state.dave; + if (!dave?.session) return; + let decrypted: Buffer; + try { + decrypted = dave.session.decrypt(uid, MediaType.VIDEO, rtpPayload); + } catch { + return; // per-packet decrypt failures are transient; skip + } + if (!decrypted || decrypted.length === 0) return; + + if (!watch.out) { + // First successful decrypt — open output file. + void openOutput(watchKey, watch, decrypted); + return; + } + watch.lastPacketAt = Date.now(); + let nals: Buffer[]; + try { + nals = watch.depacketizer.push(decrypted); + } catch { + return; + } + for (const nal of nals) { + watch.bytesWritten += nal.length; + if (!watch.out.write(nal)) { + watch.out.once("drain", () => {}); + } + } +} + +async function openOutput( + _watchKey: string, + watch: WatchState, + _first: Buffer, +): Promise { + if (watch.out) return; + const dir = path.join(_recordingsDir, watch.uid); + await fsPromises.mkdir(dir, { recursive: true }); + const filePath = path.join( + dir, + `video-${watch.channelId}-${watch.startedAt}.h264`, + ); + watch.filePath = filePath; + watch.out = createWriteStream(filePath); + logger.info( + { userId: watch.uid, path: filePath }, + "Video burst opened (DAVE video RTP received)", + ); +} + +function closeWatch(watch: WatchState): void { + if (!watch) return; + try { + watch.net.destroy(); + } catch { + /* ignore */ + } + if (watch.out && !watch.out.destroyed) { + const flushed = new Promise((resolve) => { + watch.out.once("finish", () => resolve()); + watch.out.end(); + }); + void flushed + .then(() => muxToMp4(watch.filePath)) + .then((mp4) => + logger.info( + { userId: watch.uid, mp4, bytes: watch.bytesWritten }, + "Video muxed to mp4", + ), + ) + .catch((err) => + logger.warn( + { + userId: watch.uid, + err: err instanceof Error ? err.message : String(err), + }, + "Video mux failed (raw .h264 kept)", + ), + ); + } +} + +export function closeWatchByKey(watchKey: string): void { + const watch = watches.get(watchKey); + if (!watch) return; + watches.delete(watchKey); + pendingServerId.delete(watchKey); + closeWatch(watch); +} + +/** Stop watching `userId` in `guildId` (streaming stopped / user left). */ +export function stopStreamWatch(guildId: string, userId: string): void { + closeWatchByKey(`${guildId}:${userId}`); +} + +/** Stop ALL watches for a guild (voice disconnect / channel untrack). */ +export function stopAllStreamWatches(guildId: string): void { + for (const key of [...watches.keys()]) { + if (key.startsWith(`${guildId}:`)) closeWatchByKey(key); + } +} diff --git a/services/discord-gateway/src/modules/voice-recording/videoRecorder.ts b/services/discord-gateway/src/modules/voice-recording/videoRecorder.ts index d80729fe..bb94b49b 100644 --- a/services/discord-gateway/src/modules/voice-recording/videoRecorder.ts +++ b/services/discord-gateway/src/modules/voice-recording/videoRecorder.ts @@ -1,92 +1,70 @@ /** * Video recorder (camera + screen share) for OTHER members. * - * WHY THIS EXISTS (root cause, 2026-08-31): - * Phase A/B hooked @discordjs/voice's UDP socket and depacketized H264 RTP - * directly. It captured ZERO video because @discordjs/voice is audio-only and - * NEVER sends the gateway `STREAM_WATCH` signal — so Discord never forwards a - * member's video RTP to the bot. The raw-UDP approach can't work for receive. + * ARCHITECTURE (2026-08-31, Phase D — DAVE-capable stream-watch): * - * The correct, native path lives in discord.js-selfbot-v13's own voice stack: - * - Join a selfbot `VoiceConnection` that shares the bot's single voice - * session (ClientVoiceManager.onVoiceStateUpdate feeds BOTH the - * @discordjs/voice adapter and the selfbot connection from the same - * VOICE_STATE_UPDATE — the two stacks are designed to coexist). - * - Detect a member started streaming via voiceState.streaming - * (= data.self_stream). - * - `joinStreamConnection(userId)` → sends STREAM_WATCH (op 20) so Discord - * authorizes video RTP to us. - * - `receiver.createVideoStream(userId, filepath)` → PacketHandler routes the - * member's H264+Opus RTP to a `Recorder` (ffmpeg over UDP) which muxes to - * Matroska (.mkv). + * The previous selfbot-v13 path (discord.js-selfbot-v13 `joinStreamConnection`) + * is DEAD because Discord now mandates DAVE (E2EE) on all voice RTC connections. + * The selfbot's voice stack identifies with NO `max_dave_protocol_version` → + * Discord closes the WS with code 4017 ("E2EE/DAVE protocol required") on + * every attempt, making selfbot-based video-receive impossible. * - * FAILURE FIX (2026-08-31): the selfbot `VoiceConnection` must be established - * EAGERLY at voice-join time. Calling `client.voice.joinChannel()` lazily — - * only when a member starts streaming — times out with VOICE_CONNECTION_TIMEOUT - * because the bot is already in the channel, so Discord never emits a fresh - * VOICE_SERVER_UPDATE and the Selfbot connection never gets token/endpoint. - * Establishing it alongside the @discordjs/voice join (on a FRESH join) makes - * Discord emit VOICE_SERVER_UPDATE → the selfbot connection authenticates. + * The working path: + * 1. `voiceStateUpdate` detects `streaming` on a member in a tracked channel. + * 2. `streamWatchReceiver.startStreamWatch(channel, userId)` is called. + * - It sends `STREAM_WATCH` (op 20) via the gateway WS. + * - Discord replies with `STREAM_CREATE` + `STREAM_SERVER_UPDATE` containing + * the stream's RTC endpoint/token/serverId. + * - `streamWatchReceiver` builds a `@discordjs/voice` `Networking` connection + * to the stream RTC with DAVE (the same class that handles guild audio). + * - It hooks the UDP `message` event, decrypts H264 video via Davey + * `MediaType.VIDEO`, and writes to `.h264` → muxed to `.mp4`. + * 3. `streamWatchReceiver.stopStreamWatch(guildId, userId)` tears everything down. * - * All of this is best-effort: any failure is logged and NEVER breaks the - * gateway's existing audio/voice recording. + * This module is the thin orchestration layer: detect → delegate → teardown. + * All failures are best-effort and NEVER break the gateway's audio recording. */ -import { promises as fsPromises } from "node:fs"; -import path from "node:path"; import type { Client, VoiceChannel, VoiceState } from "discord.js-selfbot-v13"; import { createChildLogger } from "@/shared/logger/index"; +import { + setStreamWatchClient, + setStreamWatchRecordingsDir, + startStreamWatch, + stopAllStreamWatches, + stopStreamWatch, +} from "./streamWatchReceiver.js"; const logger = createChildLogger("video-recorder"); -interface WatchHandle { - guildId: string; - userId: string; - recorder: unknown; - path: string; -} - -/** A connected (or in-progress) selfbot voice connection, keyed by guildId. */ -interface SelfbotVoice { - guildId: string; - conn: unknown; // selfbot VoiceConnection - connectedAt: number; -} - -/** Recorders keyed by `${guildId}:${userId}`. */ -const activeRecorders = new Map(); - /** Channels currently under audio recording, keyed by guildId. */ const watchedChannels = new Map(); -/** Eagerly-established selfbot voice connection per guild (video receive). */ -const selfbotVoices = new Map(); - let _client: Client | undefined; let _listenerAttached = false; - /** Resolved at runtime from config (kept default as fallback). */ let _recordingsDir = "/var/lib/gmw/recordings"; export function setVideoRecordingsDir(dir: string) { _recordingsDir = dir; + setStreamWatchRecordingsDir(dir); } -/** Test-only: clear all module-level state (recorders, channels, selfbot conns). */ +/** Test-only: clear all module-level state. */ export function __resetVideoRecorderState(): void { - activeRecorders.clear(); watchedChannels.clear(); - selfbotVoices.clear(); _listenerAttached = false; _client = undefined; } /** * Set the selfbot client and register the singleton voiceStateUpdate listener - * (idempotent). Called once at gateway bootstrap. + * (idempotent). Called once at gateway bootstrap. Also passes the client to + * streamWatchReceiver for raw event handling (STREAM_CREATE/SERVER_UPDATE). */ export function setVideoRecorderClient(client: Client | undefined) { _client = client; + setStreamWatchClient(client); if (!client || _listenerAttached) return; _listenerAttached = true; client.on( @@ -102,21 +80,19 @@ export function trackChannel(guildId: string, channel: VoiceChannel) { watchedChannels.set(guildId, channel); } -/** Unregister a channel (voice stopped). Tear down any video recorders. */ +/** Unregister a channel (voice stopped). Tear down any video watches. */ export function untrackChannel(guildId: string) { watchedChannels.delete(guildId); - stopAllVideoRecordings(guildId); - destroyGuildSelfbotVoice(guildId); + stopAllStreamWatches(guildId); } async function handleVoiceStateUpdate( - oldState: VoiceState, + _oldState: VoiceState, newState: VoiceState, ): Promise { try { const selfId = _client?.user?.id; if (selfId && newState.id === selfId) return; // never record bot's own video - if (_client && oldState.id === selfId) return; const guildId = newState.guild?.id; const channel = guildId ? watchedChannels.get(guildId) : undefined; @@ -126,8 +102,8 @@ async function handleVoiceStateUpdate( const inChannel = newState.channelId === channel.id; if (streaming && inChannel) { - await startVideoRecording(channel, newState.id); - } else if (!inChannel) { + startVideoRecording(channel, newState.id); + } else if (!inChannel || !streaming) { stopVideoRecording(channel.guild.id, newState.id); } } catch (err) { @@ -139,204 +115,17 @@ async function handleVoiceStateUpdate( } /** - * Resolve (or create) the selfbot `VoiceConnection` for a guild, best-effort. - * - * Establishes `client.voice.joinChannel(channel)` — the selfbot lib's own - * VoiceConnection, which has `.receiver.createVideoStream` and - * `.joinStreamConnection`. Cached per guild. Returns the connection's - * `{ conn, status }` or null on failure. - * - * IMPORTANT ordering: this MUST run while the bot is freshly joining the - * channel (alongside @discordjs/voice's join), so Discord emits VOICE_SERVER_UPDATE - * and the selfbot connection authenticates. Do NOT call it lazily after the bot - * is already connected — that times out (VOICE_CONNECTION_TIMEOUT). + * Begin watching the video (camera / screen share) of `userId` in `channel`. + * Delegates to `streamWatchReceiver` which handles STREAM_WATCH + DAVE connection. */ -export async function ensureSelfbotVoice( - channel: VoiceChannel, -): Promise<{ conn: unknown; status: number } | null> { - try { - if (!_client) { - logger.debug("ensureSelfbotVoice: no client set"); - return null; - } - const guildId = channel.guild?.id; - if (!guildId) return null; - - // Reuse an already-connected selfbot connection. - const existing = selfbotVoices.get(guildId); - if (existing && getVoiceStatus(existing.conn) === 0 /* CONNECTED */) { - return { conn: existing.conn, status: 0 }; - } - - const voiceManager = (_client as unknown as { voice?: unknown }).voice as - | { - joinChannel?: ( - ch: unknown, - cfg?: { - selfMute?: boolean; - selfDeaf?: boolean; - selfVideo?: boolean; - }, - ) => Promise; - } - | undefined; - - if (!voiceManager?.joinChannel) { - logger.warn( - { guildId }, - "ensureSelfbotVoice: selfbot voice manager unavailable", - ); - return null; - } - - logger.info( - { guildId, channelId: channel.id }, - "Establishing selfbot voice connection (video receive)", - ); - const conn = await voiceManager.joinChannel(channel, { - selfMute: false, - selfDeaf: false, - selfVideo: false, - }); - if (!conn) return null; - - const status = getVoiceStatus(conn); - selfbotVoices.set(guildId, { guildId, conn, connectedAt: Date.now() }); - logger.info( - { guildId, channelId: channel.id, status }, - status === 0 - ? "Selfbot voice connected (video receive ready)" - : "Selfbot voice connection pending", - ); - return { conn, status }; - } catch (err) { - logger.warn( - { - guildId: channel.guild?.id, - err: err instanceof Error ? err.message : String(err), - }, - "ensureSelfbotVoice failed (best-effort, ignoring)", - ); - return null; - } -} - -/** Tear down the cached selfbot voice connection for a guild. */ -export function destroyGuildSelfbotVoice(guildId: string): void { - const entry = selfbotVoices.get(guildId); - if (!entry) return; - selfbotVoices.delete(guildId); - try { - const conn = entry.conn as unknown as { - disconnect?: () => void; - destroy?: () => void; - }; - conn.disconnect?.(); - logger.info({ guildId }, "Destroyed selfbot voice connection"); - } catch (err) { - logger.warn( - { guildId, err: err instanceof Error ? err.message : String(err) }, - "Error destroying selfbot voice connection", - ); - } -} - -/** Read the VoiceStatus number off a selfbot VoiceConnection (0 = CONNECTED). */ -function getVoiceStatus(conn: unknown): number { - const status = (conn as { status?: number } | null)?.status; - return typeof status === "number" ? status : -1; -} - -/** - * Begin recording the video (camera / screen share) of `userId` in `channel`. - * Returns the watch handle on success, or null on any failure. - */ -export async function startVideoRecording( +export function startVideoRecording( channel: VoiceChannel, userId: string, -): Promise { +): void { try { - if (!_client) { - logger.warn("Video recorder: no client set, skipping"); - return null; - } - const selfId = _client.user?.id; - if (selfId && userId === selfId) return null; - - const key = `${channel.guild.id}:${userId}`; - if (activeRecorders.has(key)) return activeRecorders.get(key) ?? null; - - // Use the eagerly-established selfbot connection (falls back to a lazy - // joinChannel as a last resort — may time out if the bot already joined). - const eager = await ensureSelfbotVoice(channel); - const voiceConn: unknown | null = eager?.conn ?? null; - if (!voiceConn) return null; - - // 2. STREAM_WATCH handshake. - const watchConn = await ( - voiceConn as unknown as { - joinStreamConnection?: (u: string) => Promise; - } - ).joinStreamConnection?.(userId); - if (!watchConn) { - logger.warn( - { userId, guildId: channel.guild.id }, - "joinStreamConnection returned no connection — is the user streaming?", - ); - return null; - } - await ( - watchConn as unknown as { sendSignalScreenshare?: () => unknown } - ).sendSignalScreenshare?.(); - - // 3. Recorder → .mkv - const dir = path.join(_recordingsDir, userId); - await fsPromises.mkdir(dir, { recursive: true }); - const outPath = path.join(dir, `video-${channel.id}-${Date.now()}.mkv`); - - const receiver = (voiceConn as unknown as { receiver?: unknown }) - .receiver as - | { createVideoStream?: (u: string, out: string) => unknown } - | undefined; - if (!receiver?.createVideoStream) { - logger.warn("Video recorder: receiver.createVideoStream unavailable"); - return null; - } - const recorder = receiver.createVideoStream(userId, outPath); - - const handle: WatchHandle = { - guildId: channel.guild.id, - userId, - recorder, - path: outPath, - }; - - activeRecorders.set(key, handle); - - const rec = recorder as unknown as { - on?: (e: string, cb: () => void) => void; - }; - if (typeof rec.on === "function") { - rec.on("ready", () => { - logger.info( - { userId, path: outPath, guildId: channel.guild.id }, - "Video recorder ready (ffmpeg muxing started)", - ); - }); - rec.on("closed", () => { - logger.info( - { userId, path: outPath, guildId: channel.guild.id }, - "Video recorder closed", - ); - activeRecorders.delete(key); - }); - } - - logger.info( - { userId, path: outPath, guildId: channel.guild.id }, - "Started video recording via STREAM_WATCH", - ); - return handle; + const selfId = _client?.user?.id; + if (selfId && userId === selfId) return; + startStreamWatch(channel, userId); } catch (err) { logger.warn( { @@ -344,41 +133,17 @@ export async function startVideoRecording( guildId: channel.guild.id, err: err instanceof Error ? err.message : String(err), }, - "Video recording failed (best-effort, ignoring)", + "startVideoRecording error (best-effort, ignoring)", ); - return null; } } -/** Stop and tear down the video recorder for userId in guildId. */ +/** Stop watching `userId`'s video in `guildId`. */ export function stopVideoRecording(guildId: string, userId: string): void { - const key = `${guildId}:${userId}`; - const handle = activeRecorders.get(key); - if (!handle) return; - try { - const rec = handle.recorder as unknown as { - destroy?: () => void; - close?: () => void; - }; - rec.destroy?.(); - activeRecorders.delete(key); - logger.info({ userId, path: handle.path }, "Stopped video recording"); - } catch (err) { - logger.warn( - { userId, err: err instanceof Error ? err.message : String(err) }, - "Error stopping video recording", - ); - } + stopStreamWatch(guildId, userId); } -/** Tear down ALL video recorders for a guild (e.g. on voice disconnect). */ +/** Tear down ALL video watches for a guild (e.g. on voice disconnect). */ export function stopAllVideoRecordings(guildId: string): void { - for (const [key, handle] of activeRecorders) { - if (handle?.guildId === guildId) { - stopVideoRecording(handle.guildId, handle.userId); - activeRecorders.delete(key); - } - } + stopAllStreamWatches(guildId); } - -export type { WatchHandle }; diff --git a/services/discord-gateway/tests/videoRecorder.test.ts b/services/discord-gateway/tests/videoRecorder.test.ts index 0ae12f87..6ec0d3ce 100644 --- a/services/discord-gateway/tests/videoRecorder.test.ts +++ b/services/discord-gateway/tests/videoRecorder.test.ts @@ -1,8 +1,6 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; import { __resetVideoRecorderState, - destroyGuildSelfbotVoice, - ensureSelfbotVoice, setVideoRecorderClient, setVideoRecordingsDir, startVideoRecording, @@ -10,26 +8,9 @@ import { trackChannel, untrackChannel, } from "../src/modules/voice-recording/videoRecorder.js"; +import * as streamWatch from "../src/modules/voice-recording/streamWatchReceiver.js"; -// ─── mock the selfbot VoiceConnection/watch/recorder surface ───────────── -function makeVoiceManager() { - const recorder = { - on: vi.fn(), - destroy: vi.fn(), - }; - const watchConn = { sendSignalScreenshare: vi.fn(async () => {}) }; - const voiceConn = { - status: 0, // VoiceStatus.CONNECTED - disconnect: vi.fn(), - receiver: { - createVideoStream: vi.fn(() => recorder), - }, - joinStreamConnection: vi.fn(async () => watchConn), - }; - const joinChannel = vi.fn(async () => voiceConn); - return { voiceConn, watchConn, recorder, joinChannel }; -} - +// ─── mocks ───────────────────────────────────────────────────────────── function makeChannel(guildId = "g1", channelId = "c1") { return { id: channelId, @@ -40,7 +21,6 @@ function makeChannel(guildId = "g1", channelId = "c1") { function makeClient() { return { user: { id: "bot1" }, - voice: makeVoiceManager(), on: vi.fn(), }; } @@ -53,109 +33,72 @@ function makeUser() { beforeEach(() => { __resetVideoRecorderState(); + vi.restoreAllMocks(); setVideoRecordingsDir("/tmp/gmw-vidrec-test"); }); -describe("videoRecorder", () => { - it("registers a single voiceStateUpdate listener on the client", () => { +describe("videoRecorder (stream-watch orchestration)", () => { + it("registers a voiceStateUpdate + raw listener on the client (idempotent)", () => { const client = makeClient(); setVideoRecorderClient(client); - setVideoRecorderClient(client); // idempotent - expect(client.on).toHaveBeenCalledTimes(1); + setVideoRecorderClient(client); // idempotent (videoRecorder + streamWatch) + // one raw (streamWatch) + one voiceStateUpdate (videoRecorder) + expect(client.on).toHaveBeenCalledTimes(2); expect(client.on).toHaveBeenCalledWith( "voiceStateUpdate", expect.any(Function), ); + expect(client.on).toHaveBeenCalledWith("raw", expect.any(Function)); }); - it("startVideoRecording does watch handshake + createVideoStream to a .mkv path", async () => { - const client = makeClient(); - const { watchConn } = client.voice; + it("startVideoRecording delegates to streamWatchReceiver.startStreamWatch", () => { + const spy = vi.spyOn(streamWatch, "startStreamWatch").mockImplementation(() => {}); + const ch = makeChannel(); + startVideoRecording(ch, makeUser()); + expect(spy).toHaveBeenCalledWith(ch, expect.any(String)); + }); + + it("refuses to record the bot's own video", () => { + setVideoRecorderClient(makeClient()); + const spy = vi.spyOn(streamWatch, "startStreamWatch").mockImplementation(() => {}); + startVideoRecording(makeChannel(), "bot1"); + expect(spy).not.toHaveBeenCalled(); + }); + + it("stopVideoRecording calls streamWatchReceiver.stopStreamWatch", () => { + const spy = vi.spyOn(streamWatch, "stopStreamWatch").mockImplementation(() => {}); + const u = makeUser(); + stopVideoRecording("g1", u); + expect(spy).toHaveBeenCalledWith("g1", u); + }); + + it("untrackChannel tears down all stream watches for the guild", () => { + const spy = vi.spyOn(streamWatch, "stopAllStreamWatches").mockImplementation(() => {}); + untrackChannel("g1"); + expect(spy).toHaveBeenCalledWith("g1"); + }); + + it("the voiceStateUpdate handler starts a watch when a member starts streaming", () => { + const client = makeClient() as any; setVideoRecorderClient(client); trackChannel("g1", makeChannel()); - - const handle = await startVideoRecording(makeChannel(), makeUser()); - - expect(handle).not.toBeNull(); - expect(handle?.path).toMatch(/video-c1-\d+\.mkv$/); - expect(client.voice.voiceConn.joinStreamConnection).toHaveBeenCalled(); - expect(watchConn.sendSignalScreenshare).toHaveBeenCalled(); - expect( - client.voice.voiceConn.receiver.createVideoStream, - ).toHaveBeenCalledWith(expect.any(String), expect.stringMatching(/\.mkv$/)); + const startSpy = vi.spyOn(streamWatch, "startStreamWatch").mockImplementation(() => {}); + const listener = client.on.mock.calls.find( + (c: unknown[]) => c[0] === "voiceStateUpdate", + )?.[1]; + listener(null, { id: "u1", guild: { id: "g1" }, channelId: "c1", streaming: true }); + expect(startSpy).toHaveBeenCalled(); }); - it("refuses to record the bot's own video", async () => { - const client = makeClient(); + it("the voiceStateUpdate handler stops a watch when a member stops streaming", () => { + const client = makeClient() as any; setVideoRecorderClient(client); - const handle = await startVideoRecording(makeChannel(), "bot1"); - expect(handle).toBeNull(); - expect(client.voice.voiceConn.joinStreamConnection).not.toHaveBeenCalled(); - }); - - it("is idempotent for the same guild:user (no duplicate recorder)", async () => { - const client = makeClient(); - setVideoRecorderClient(client); - const ch = makeChannel(); - const u = makeUser(); - await startVideoRecording(ch, u); - await startVideoRecording(ch, u); - expect( - client.voice.voiceConn.receiver.createVideoStream, - ).toHaveBeenCalledTimes(1); - }); - - it("stopVideoRecording destroys the recorder", async () => { - const client = makeClient(); - setVideoRecorderClient(client); - const ch = makeChannel(); - const u = makeUser(); - await startVideoRecording(ch, u); - stopVideoRecording("g1", u); - expect(client.voice.recorder.destroy).toHaveBeenCalled(); - }); - - it("untrackChannel tears down active video recorders for the guild", async () => { - const client = makeClient(); - setVideoRecorderClient(client); - const ch = makeChannel(); - const u = makeUser(); - await startVideoRecording(ch, u); - untrackChannel("g1"); - expect(client.voice.recorder.destroy).toHaveBeenCalled(); - }); - - it("ensureSelfbotVoice establishes and caches the selfbot connection (join once)", async () => { - const client = makeClient(); - setVideoRecorderClient(client); - const ch = makeChannel(); - const r1 = await ensureSelfbotVoice(ch); - expect(r1?.status).toBe(0); - expect(client.voice.joinChannel).toHaveBeenCalledTimes(1); - // Second call reuses the cached CONNECTED connection — no re-join. - const r2 = await ensureSelfbotVoice(makeChannel()); - expect(client.voice.joinChannel).toHaveBeenCalledTimes(1); - expect(r2?.status).toBe(0); - }); - - it("untrackChannel destroys the cached selfbot voice connection", async () => { - const client = makeClient(); - setVideoRecorderClient(client); - await ensureSelfbotVoice(makeChannel("g1", "c1")); - const other = makeClient(); - setVideoRecorderClient(other); - await ensureSelfbotVoice(makeChannel("g2", "c2")); - untrackChannel("g1"); - expect(client.voice.voiceConn.disconnect).toHaveBeenCalled(); - // g2 connection untouched by g1 teardown - expect(other.voice.voiceConn.disconnect).not.toHaveBeenCalled(); - }); - - it("destroyGuildSelfbotVoice disconnects the selfbot connection", async () => { - const client = makeClient(); - setVideoRecorderClient(client); - await ensureSelfbotVoice(makeChannel()); - destroyGuildSelfbotVoice("g1"); - expect(client.voice.voiceConn.disconnect).toHaveBeenCalled(); + trackChannel("g1", makeChannel()); + const stopSpy = vi.spyOn(streamWatch, "stopStreamWatch").mockImplementation(() => {}); + const listener = client.on.mock.calls.find( + (c: unknown[]) => c[0] === "voiceStateUpdate", + )?.[1]; + listener(null, { id: "u1", guild: { id: "g1" }, channelId: "c1", streaming: false }); + expect(stopSpy).toHaveBeenCalledWith("g1", "u1"); }); });