feat(gateway): record others' video via native selfbot watch/receive (Phase C)
Root cause of why Phase A/B captured zero video: @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. Live diagnostic confirmed: while members were sharing, audio .ogg files flowed for many users but no non-opus RTP ever arrived. Correct path: discord.js-selfbot-v13 ships a complete native watch/record stack. New src/modules/voice-recording/videoRecorder.ts: - Detects streamers via voiceState.streaming on a single idempotent voiceStateUpdate listener. - client.voice.joinChannel() (reuses the single session alongside @discordjs/voice) + joinStreamConnection(userId) -> STREAM_WATCH (op 20). - receiver.createVideoStream(userId, path) -> PacketHandler -> Recorder (ffmpeg over UDP loopback) -> Matroska .mkv, decryption handled internally. Wired in recorder.ts (trackChannel/untrackChannel) + bootstrap.ts (setVideoRecorderClient/setVideoRecordingsDir). All best-effort; failures log and never break existing voice/audio. Unit tests 6/6 (single listener, watch handshake + mkv path, skip own video, idempotence, teardown). Full suite 20 files / 176 tests green; typecheck + build + biome clean. UNVERIFIED live yet: needs deploy + a streamer to confirm Recorder ready + playable .mkv.
This commit is contained in:
@@ -36,6 +36,10 @@ import {
|
||||
setPcmWsClient,
|
||||
setEventBroadcaster as setRecorderEventBroadcaster,
|
||||
} from "../modules/voice-recording/recorder.js";
|
||||
import {
|
||||
setVideoRecorderClient,
|
||||
setVideoRecordingsDir,
|
||||
} from "../modules/voice-recording/videoRecorder.js";
|
||||
import { VoiceController } from "../modules/voice-recording/voiceController.js";
|
||||
import { config } from "../shared/config/config.js";
|
||||
import {
|
||||
@@ -182,6 +186,12 @@ export async function initializeDiscordGateway() {
|
||||
const client = new Client(createDiscordClientOptions());
|
||||
const voiceController = new VoiceController(client);
|
||||
|
||||
// Wire the video recorder (others' camera/screen share) to the selfbot's
|
||||
// native watch/receive stack + its recordings dir. Best-effort: failures are
|
||||
// logged inside, never fatal.
|
||||
setVideoRecorderClient(client);
|
||||
setVideoRecordingsDir(config.RECORDINGS_DIR);
|
||||
|
||||
// Initialize Redis event broadcaster
|
||||
const redisPublisher = new RedisEventPublisher(config.REDIS_URL, logger);
|
||||
const eventBroadcaster = new EventBroadcaster(redisPublisher);
|
||||
|
||||
@@ -20,6 +20,7 @@ import {
|
||||
import { createSpeakingHandler } from "./recorder/speakingHandler.js";
|
||||
import { hookScreenShareAudio } from "./screenShareAudio.js";
|
||||
import { hookVideoReceiver } from "./videoReceiver.js";
|
||||
import { trackChannel, untrackChannel } from "./videoRecorder.js";
|
||||
|
||||
const logger = createChildLogger("recorder");
|
||||
|
||||
@@ -163,6 +164,8 @@ export async function startRecording(
|
||||
startTime: sessionStartTime,
|
||||
});
|
||||
activeSessions.set(channel.guild.id, session);
|
||||
// Register this channel for video recording (others' camera/screen share).
|
||||
trackChannel(channel.guild.id, channel);
|
||||
} catch (err) {
|
||||
logger.error({ error: err }, "Failed to connect to voice channel");
|
||||
connection.destroy();
|
||||
@@ -230,6 +233,7 @@ export async function startRecording(
|
||||
|
||||
connection.on(VoiceConnectionStatus.Destroyed, () => {
|
||||
finalizeActiveRecordingSession(channel.guild.id);
|
||||
untrackChannel(channel.guild.id);
|
||||
if (config.VERBOSE) {
|
||||
logger.info("Voice connection destroyed");
|
||||
}
|
||||
@@ -251,6 +255,7 @@ export function stopRecording(guildId: string): void {
|
||||
} else {
|
||||
logger.warn("No active connection to stop");
|
||||
}
|
||||
untrackChannel(guildId);
|
||||
|
||||
const session = activeSessions.get(guildId);
|
||||
finalizeActiveRecordingSession(guildId);
|
||||
|
||||
@@ -0,0 +1,268 @@
|
||||
/**
|
||||
* 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.
|
||||
*
|
||||
* 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).
|
||||
*
|
||||
* All of this is best-effort: any failure is logged and NEVER breaks the
|
||||
* gateway's existing audio/voice 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";
|
||||
|
||||
const logger = createChildLogger("video-recorder");
|
||||
|
||||
interface WatchHandle {
|
||||
guildId: string;
|
||||
userId: string;
|
||||
recorder: unknown;
|
||||
path: string;
|
||||
}
|
||||
|
||||
/** Recorders keyed by `${guildId}:${userId}`. */
|
||||
const activeRecorders = new Map<string, WatchHandle>();
|
||||
|
||||
/** Channels currently under audio recording, keyed by guildId. */
|
||||
const watchedChannels = new Map<string, VoiceChannel>();
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the selfbot client and register the singleton voiceStateUpdate listener
|
||||
* (idempotent). Called once at gateway bootstrap.
|
||||
*/
|
||||
export function setVideoRecorderClient(client: Client | undefined) {
|
||||
_client = client;
|
||||
if (!client || _listenerAttached) return;
|
||||
_listenerAttached = true;
|
||||
client.on(
|
||||
"voiceStateUpdate",
|
||||
(oldState: VoiceState, newState: VoiceState) => {
|
||||
void handleVoiceStateUpdate(oldState, newState);
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
/** Register a channel whose audio recording is active (video should follow). */
|
||||
export function trackChannel(guildId: string, channel: VoiceChannel) {
|
||||
watchedChannels.set(guildId, channel);
|
||||
}
|
||||
|
||||
/** Unregister a channel (voice stopped). Tear down any video recorders. */
|
||||
export function untrackChannel(guildId: string) {
|
||||
watchedChannels.delete(guildId);
|
||||
stopAllVideoRecordings(guildId);
|
||||
}
|
||||
|
||||
async function handleVoiceStateUpdate(
|
||||
oldState: VoiceState,
|
||||
newState: VoiceState,
|
||||
): Promise<void> {
|
||||
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;
|
||||
if (!channel) return; // not an actively-recorded channel
|
||||
|
||||
const streaming = Boolean(newState.streaming);
|
||||
const inChannel = newState.channelId === channel.id;
|
||||
|
||||
if (streaming && inChannel) {
|
||||
await startVideoRecording(channel, newState.id);
|
||||
} else if (!inChannel) {
|
||||
stopVideoRecording(channel.guild.id, newState.id);
|
||||
}
|
||||
} catch (err) {
|
||||
logger.warn(
|
||||
{ err: err instanceof Error ? err.message : String(err) },
|
||||
"Video voice-state handler error (ignoring)",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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(
|
||||
channel: VoiceChannel,
|
||||
userId: string,
|
||||
): Promise<WatchHandle | null> {
|
||||
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;
|
||||
|
||||
const voiceManager = (_client as unknown as { voice?: unknown }).voice as
|
||||
| {
|
||||
joinChannel?: (
|
||||
ch: unknown,
|
||||
cfg?: {
|
||||
selfMute?: boolean;
|
||||
selfDeaf?: boolean;
|
||||
selfVideo?: boolean;
|
||||
},
|
||||
) => Promise<unknown>;
|
||||
}
|
||||
| undefined;
|
||||
|
||||
if (!voiceManager?.joinChannel) {
|
||||
logger.warn("Video recorder: selfbot voice manager unavailable");
|
||||
return null;
|
||||
}
|
||||
|
||||
// 1. Selfbot voice connection — joinChannel is idempotent (reuses
|
||||
// ClientVoiceManager.connection, re-confirms the same session).
|
||||
const voiceConn = await voiceManager.joinChannel(channel, {
|
||||
selfMute: false,
|
||||
selfDeaf: false,
|
||||
selfVideo: false,
|
||||
});
|
||||
if (!voiceConn) return null;
|
||||
|
||||
// 2. STREAM_WATCH handshake.
|
||||
const watchConn = await (
|
||||
voiceConn as unknown as {
|
||||
joinStreamConnection?: (u: string) => Promise<unknown>;
|
||||
}
|
||||
).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;
|
||||
} catch (err) {
|
||||
logger.warn(
|
||||
{
|
||||
userId,
|
||||
guildId: channel.guild.id,
|
||||
err: err instanceof Error ? err.message : String(err),
|
||||
},
|
||||
"Video recording failed (best-effort, ignoring)",
|
||||
);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** Stop and tear down the video recorder for userId 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",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/** Tear down ALL video recorders 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export type { WatchHandle };
|
||||
@@ -0,0 +1,120 @@
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
setVideoRecorderClient,
|
||||
setVideoRecordingsDir,
|
||||
startVideoRecording,
|
||||
stopVideoRecording,
|
||||
trackChannel,
|
||||
untrackChannel,
|
||||
} from "../src/modules/voice-recording/videoRecorder.js";
|
||||
|
||||
// ─── mock the selfbot VoiceConnection/watch/recorder surface ─────────────
|
||||
function makeVoiceManager() {
|
||||
const recorder = {
|
||||
on: vi.fn(),
|
||||
destroy: vi.fn(),
|
||||
};
|
||||
const voiceConn = {
|
||||
receiver: {
|
||||
createVideoStream: vi.fn(() => recorder),
|
||||
},
|
||||
joinStreamConnection: vi.fn(async () => watchConn),
|
||||
};
|
||||
const watchConn = { sendSignalScreenshare: vi.fn(async () => {}) };
|
||||
return { voiceConn, watchConn, recorder, joinChannel: async () => voiceConn };
|
||||
}
|
||||
|
||||
function makeChannel(guildId = "g1", channelId = "c1") {
|
||||
return {
|
||||
id: channelId,
|
||||
guild: { id: guildId },
|
||||
};
|
||||
}
|
||||
|
||||
function makeClient() {
|
||||
return {
|
||||
user: { id: "bot1" },
|
||||
voice: makeVoiceManager(),
|
||||
on: vi.fn(),
|
||||
};
|
||||
}
|
||||
|
||||
let seq = 0;
|
||||
function makeUser() {
|
||||
seq += 1;
|
||||
return `user-${seq}`;
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
setVideoRecordingsDir("/tmp/gmw-vidrec-test");
|
||||
});
|
||||
|
||||
describe("videoRecorder", () => {
|
||||
it("registers a single voiceStateUpdate listener on the client", () => {
|
||||
const client = makeClient();
|
||||
setVideoRecorderClient(client);
|
||||
setVideoRecorderClient(client); // idempotent
|
||||
expect(client.on).toHaveBeenCalledTimes(1);
|
||||
expect(client.on).toHaveBeenCalledWith(
|
||||
"voiceStateUpdate",
|
||||
expect.any(Function),
|
||||
);
|
||||
});
|
||||
|
||||
it("startVideoRecording does watch handshake + createVideoStream to a .mkv path", async () => {
|
||||
const client = makeClient();
|
||||
const { watchConn } = client.voice;
|
||||
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$/));
|
||||
});
|
||||
|
||||
it("refuses to record the bot's own video", async () => {
|
||||
const client = makeClient();
|
||||
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();
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user