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.
This commit is contained in:
@@ -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<VoiceConnection | null> {
|
||||
// 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,
|
||||
|
||||
@@ -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:<gid>:<chid>:<uid>`; 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<string, WatchState>(); // key `${guildId}:${uid}`
|
||||
const pendingServerId = new Map<string, string>(); // 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:<gid>:<chid>:<uid>` 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<string, unknown> } | 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<void> {
|
||||
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<void>((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);
|
||||
}
|
||||
}
|
||||
@@ -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<string, WatchHandle>();
|
||||
|
||||
/** Channels currently under audio recording, keyed by guildId. */
|
||||
const watchedChannels = new Map<string, VoiceChannel>();
|
||||
|
||||
/** Eagerly-established selfbot voice connection per guild (video receive). */
|
||||
const selfbotVoices = new Map<string, SelfbotVoice>();
|
||||
|
||||
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<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;
|
||||
@@ -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<unknown>;
|
||||
}
|
||||
| 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<WatchHandle | null> {
|
||||
): 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<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;
|
||||
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 };
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user