Compare commits

12 Commits
Author SHA1 Message Date
asepharyana f70a92880e feat(glossary): implement term glossary for LLM moderation with caching and extraction logic 2026-08-12 13:44:52 +07:00
asepharyana 88b13225cd fix(gateway): stream screen-share input from yt-dlp stdout — no more raw-URL 403
Second root cause (2026-08-12): even with yt-dlp http_headers forwarded,
YouTube still returns 403 when a signed DASH URL from --dump-single-json is
fetched raw by ffmpeg/curl on some videos (verified on fONoh7Pc6VU: curl
with the EXACT headers got 403; yt-dlp's own downloader succeeded). The
signature is tied to the extracting client context (po_token/visitor), not
just UA/IP.

Fix: getDirectScreenInput now spawns 'yt-dlp -o -' and returns its stdout
as a Readable — the same mechanism resolveMediaUrl already uses for music.
yt-dlp handles auth, cookies and transient retries internally. Merge
fragments go to /tmp/gmw-ytdlp-tmp (Nix store CWD is read-only → EACCES).
Removed resolveScreenInput + mergeScreenStreams (dead code).

Controller resolveInputWithRetry unchanged: tees the stream, waits for the
first byte (12s), retries with a fresh yt-dlp run up to 3x on error/EOF/
timeout, and destroys stuck inputs (EPIPE) so no process leaks.

Tests: rewritten for streaming (yt-dlp emits bytes; fail mode = exit 8
without stdout → stream must terminate with zero bytes).
2026-08-12 13:23:59 +07:00
asepharyana 67ab289caa fix(gateway): fail-fast + retry on screen share merge failure (black tile zombie)
Root cause (2026-08-12 11:50 test): merge ffmpeg hit a transient YouTube
403 and exited code 8 BEFORE prepareStream attached its input listeners
(voice release+join takes ~10s). The input's end/error events fired into
the void, the encoder stdin never received EOF, demux resolved with
fallback 0x0 metadata, setSpeaking fired anyway → stream 'started' with
zero frames for 8+ minutes (black tile, both ffmpeg processes hung).

Fixes:
- mediaSource: pass yt-dlp http_headers (UA/referer) to the merge ffmpeg
  via -headers to suppress transient 403s; destroy the returned stream
  with an error when the merge exits non-zero before producing bytes.
- screenShareController: resolveInputWithRetry — tee the merge stream and
  wait for the first readable byte (12s timeout) before proceeding; on
  error/EOF/timeout retry the whole resolution with a FRESH yt-dlp run
  (signed DASH URLs expire fast) up to 3 attempts. Stuck merges get
  EPIPE via input.destroy() so no process leaks per attempt.
- prepareStream: race guard — if the input already ended/destroyed before
  listeners attach, EOF the encoder stdin immediately; first-frame
  watchdog in playStream rejects 'started but nothing flowing' after 10s
  instead of resolving with a silent black stream.

Tests: +2 (merge-fail zero-byte terminal state, -headers forwarding).
2026-08-12 12:17:43 +07:00
asepharyana ef9e243609 fix(goLive): encode H264 baseline to match SDP profile-level-id (black tile)
SDP offer advertises profile-level-id=42e01f (constrained baseline) but
x264 encoded the default High profile — Discord's receiver configures its
decoder from the negotiated profile, so the High-profile bitstream failed
to decode → black GoLive tile despite valid access units + correct RTP
timestamps (fixed in 42a503c).

- Add -profile:v baseline to H264 encoder options (matches @dank074's
  proven config; SPS now 6742c01e → profile_idc=66 baseline, aligns with
  the 42e01f fmtp).
- Default x264 tune film → zerolatency (no lookahead — correct for live
  GoLive; @dank074 uses it).
- Update goLive-port test to assert baseline + zerolatency.
2026-08-12 11:37:04 +07:00
asepharyana 42a503c206 fix(goLive): demux access-unit grouping + correct RTP timestamps (black tile)
Demuxer emitted each AnnexB NAL as its own WebRTC frame (SPS/PPS/SEI
separate from slices) with a near-zero timestamp delta (duration=1 in a
1/90000 timebase → RTP +1/frame instead of +3000 @30fps). Discord's H264
receiver never receives a complete decodable access unit → black GoLive
tile despite frames flowing.

- Group NALs into access units: buffer param-set/SEI NALs, flush one
  frame per slice with preceding parameter sets (AnnexB start codes kept
  so the H264RtpPacketizer finds NAL boundaries).
- Timestamp each frame at the video frame rate: duration=1, timeBase
  1/fps → BaseMediaStream frametime=1000/fps ms → RTP +clockRate/fps
  (3000 @ 30fps/90kHz) and correct pacing.
- Thread explicit frameRate from playStream options (raw H264 has no
  timing info; ffmpeg guesses 25fps on stderr).
- Strengthen golive-demux-live-e2e: validates every frame has a slice,
  no bare param-set frames, keyframes carry SPS/PPS, timeBase 1/30.
2026-08-12 11:13:00 +07:00
asepharyana 652974e23a fix(goLive): gateway crash on screenshare stop — unhandledRejection during teardown
Test 00:32 confirmed the video pipeline WORKS (1410 frames @ 1280x720 sent,
ready=true, camera off) but the gateway crashed at stream stop:
unhandledRejection → graceful shutdown → systemd restart (bot offline).

Root cause candidates (both were fire-and-forget promises without .catch):
- BaseMediaConnection.setProtocols().then(...) — rejects when the PC is
  closed while setProtocols is in flight (stream teardown)
- void webRtcConn.createOffer().then(...) — rejects when the PC closes
  while the offer is still gathering

Fixes:
- .catch on both promise chains (log + continue; teardown is expected)
- unhandledRejection handler now treats transient stream errors (EPIPE,
  ERR_STREAM_DESTROYED, ERR_STREAM_WRITE_AFTER_END, ECONNRESET) like
  uncaughtException already does — warn + continue instead of shutting
  down the whole gateway. Non-transient rejections still log + shutdown
  (with String(reason) so the detail actually shows).
2026-08-12 00:36:53 +07:00
asepharyana 407e003399 fix(goLive): black screen root cause — h264 muxer can't carry audio; disable self_video camera
ROOT CAUSE of empty GoLive tile (finally): prepareStream ran with
includeAudio: true + output -f h264. The h264 muxer cannot mux audio
('h264 muxer does not support any stream of type audio') → header write
fails -22 → stdout empty → Demuxer ffmpeg 'Invalid data found when
processing input' → 0 frames → black tile. Reproduced locally end-to-end
(13s backpressure delay + prepareStream + demux).

Fixes:
- screenShareController: includeAudio: false (video-only GoLive; demux
  path never delivers audio anyway)
- Demuxer: pin input format -f h264 for stream inputs (raw AnnexB H264
  has no magic header → auto-detect unreliable on delayed pipes)
- Streamer.signalStream: self_video: false — stop flipping on the bot's
  camera in Discord (user request; screen share ≠ camera)

Verified: local repro now emits 644 frames 1280x720 (was 0); tsc/biome/
vitest all green.
2026-08-12 00:25:01 +07:00
asepharyana 968a43b0f4 debug(goLive): instrument frame pipeline — demux spawn/stderr/frames, playStream resolve, sendVideoFrame drop/send
Tile kosong meski STREAM_CREATE handshake penuh (22:18-22:19 retest):
- Demuxer logs spawn args, ffmpeg stderr errors, frame count every 30
- playStream logs createStream resolved + demux done + setPacketizer
- sendVideoFrame logs DROPPED (ready/track) + sent frame count
2026-08-11 22:31:31 +07:00
asepharyana 91c7a67d2f fix(goLive): stream demux directly instead of spool-to-file (empty screen share)
Root cause of 'tile appears but content empty': demux() spooled the live
NUT/H264 input to a temp file and awaited stream 'finish' — but the merge
ffmpeg output never ends during playback, so demux deadlocked, no probe,
no transcode, 0 frames sent.

- Demuxer: pipe input straight into ffmpeg stdin (-i pipe:0), parse NAL
  frames live from stdout; parse video metadata from ffmpeg stderr with a
  1.5s race (fall back to H264 defaults). No spool, no await-end.
- screenShareController: pass width/height/frameRate (1280x720@30) to
  playStream — matches the prepareStream encode settings, so setVideoAttributes
  gets real dimensions even when ffmpeg can't report metadata on an open pipe.
- Add tests/golive-demux-live-e2e.ts: proves frames flow while input is
  still open (regression test for the deadlock).
2026-08-11 21:43:16 +07:00
asepharyana 8615383829 fix(goLive): STREAM_CREATE handshake — self_video voice state + retry + send instrumentation
- signalStream: flip voice state to self_video:true/self_deaf:false before
  STREAM_CREATE (Discord silently ignores the request while video disabled)
- createStream: attach dispatch listeners before first signal (race), clean
  up listeners on timeout, retry STREAM_CREATE every 3s up to 4 attempts
  (upstream issue #217/#219 — Discord randomly drops the request)
- sendOpcode: direct [goLive:Streamer] log bypassing bootstrap debug filter
  (proves op 18 is actually broadcast)
2026-08-11 20:45:00 +07:00
asepharyana 10d7ecd405 fix(gateway): EPIPE crash on media stop — stream error handlers + no shutdown on transient stream errors 2026-08-11 20:10:05 +07:00
asepharyana ff554fcff2 fix(goLive): instrument voice/stream handshake + createStream timeout (12s) 2026-08-11 20:00:33 +07:00
21 changed files with 1525 additions and 472 deletions
+40 -1
View File
@@ -279,12 +279,51 @@ export async function initializeDiscordGateway() {
});
process.on("uncaughtException", (err) => {
const code =
typeof (err as NodeJS.ErrnoException).code === "string"
? (err as NodeJS.ErrnoException).code
: "";
// Transient stream-teardown errors (voice stop/disconnect races, child
// process stdin closed while we still write) are NOT fatal — crashing the
// gateway on EPIPE takes the whole bot offline mid-music. Log + continue.
if (
code === "EPIPE" ||
code === "ERR_STREAM_DESTROYED" ||
code === "ERR_STREAM_WRITE_AFTER_END" ||
code === "ECONNRESET"
) {
logger.warn(
{ error: err },
"Uncaught transient stream error — continuing",
);
return;
}
logger.error({ error: err }, "Uncaught exception");
gracefulShutdown("uncaughtException");
});
process.on("unhandledRejection", (reason, promise) => {
logger.error({ reason, promise }, "Unhandled rejection");
const err =
reason instanceof Error ? reason : new Error(String(reason ?? "unknown"));
const code = (err as NodeJS.ErrnoException).code ?? "";
// Same transient-teardown policy as uncaughtException: a rejection that
// fires while a stream is being torn down (EPIPE after ffmpeg stdin
// closes, write-after-destroy, socket reset) must NOT take the whole
// gateway offline. Log detail + continue. Everything else still shuts
// down so real bugs surface.
if (
code === "EPIPE" ||
code === "ERR_STREAM_DESTROYED" ||
code === "ERR_STREAM_WRITE_AFTER_END" ||
code === "ECONNRESET"
) {
logger.warn(
{ error: err },
"Unhandled rejection transient stream error — continuing",
);
return;
}
logger.error({ error: err, reason: String(reason) }, "Unhandled rejection");
gracefulShutdown("unhandledRejection");
});
@@ -262,6 +262,9 @@ a=ice-lite
[audioSection, videoSection, videoRtpMap].join("\n"),
"answer",
);
console.log(
`[goLive:${this.constructor.name}] SELECT_PROTOCOL_ACK processed — remote answer set (${[audioSection, videoSection].join("\n").length}B)`,
);
this.emit("select_protocol_ack");
}
@@ -330,7 +333,15 @@ a=ice-lite
if (seq) this._sequenceNumber = seq;
if (op === VoiceOpCodes.READY) {
this.handleReady(d);
this.setProtocols().then(() => this.ready?.(this._webRtcWrapper));
this.setProtocols()
.then(() => this.ready?.(this._webRtcWrapper))
.catch((err: unknown) => {
// PC can be closed while setProtocols is in flight (stream
// teardown) — don't let that become an unhandledRejection.
console.log(
`[goLive:${this.constructor.name}] setProtocols rejected during teardown: ${err instanceof Error ? err.message : String(err)}`,
);
});
this.setVideoAttributes(false);
} else if (op >= 4000) {
console.error(`${this.constructor.name} connection error`, d);
@@ -519,10 +530,14 @@ a=ice-lite
const reconnect = () => {
const webRtcConn = this._webRtcWrapper.initWebRtc();
webRtcConn.onStateChange((state) => {
console.log(`[goLive:${this.constructor.name}] pc state => ${state}`);
if (state === "closed" && !this._closed) reconnect();
});
this._webRtcWrapper.onLocalDescription = (sdp) => {
const rtc_connection_id = randomUUID();
console.log(
`[goLive:${this.constructor.name}] sending SELECT_PROTOCOL (offer ${sdp.length}B, rtc_connection_id=${rtc_connection_id.slice(0, 8)})`,
);
this.sendOpcode(VoiceOpCodes.SELECT_PROTOCOL, {
protocol: "webrtc",
codecs: Object.values(CodecPayloadType),
@@ -532,9 +547,18 @@ a=ice-lite
});
};
// createOffer (binding resolves full SDP incl. candidates after gathering)
void webRtcConn.createOffer().then((sdp) => {
this._webRtcWrapper.onLocalDescription?.(sdp);
});
void webRtcConn
.createOffer()
.then((sdp) => {
this._webRtcWrapper.onLocalDescription?.(sdp);
})
.catch((err: unknown) => {
// PC closed while offer is gathering (stream teardown / reconnect) —
// swallow, the reconnect loop will start a fresh offer.
console.log(
`[goLive:${this.constructor.name}] createOffer rejected: ${err instanceof Error ? err.message : String(err)}`,
);
});
};
reconnect();
return new Promise((resolve) => {
+182 -123
View File
@@ -13,12 +13,13 @@
*/
import { spawn } from "node:child_process";
import { randomUUID } from "node:crypto";
import { createWriteStream, existsSync, readdirSync } from "node:fs";
import { tmpdir } from "node:os";
import { existsSync, readdirSync } from "node:fs";
import { join } from "node:path";
import { PassThrough } from "node:stream";
/** 4-byte AnnexB start code (00 00 00 01) used when building access units. */
const startCode4 = Buffer.from([0, 0, 0, 1]);
/**
* Resolve ffmpeg/ffprobe binary. Prefers explicit env override, then PATH,
* then a Nix-store ffmpeg-headless (the GMW flake provides it in the service
@@ -140,121 +141,36 @@ export async function probeStreams(
/**
* Demux input (URL string or readable stream) into video frames on a
* PassThrough. Uses ffmpeg -f h264 -c copy for video-only AnnexB output.
* Returns stream info + the video pipe. Audio is not extracted (GoLive
* screen share sends silence / uses Discord's mixed audio).
* PassThrough. Streams input DIRECTLY into ffmpeg (no spool-to-file — the
* live NUT/H264 source never ends, so spooling deadlocks). ffmpeg emits
* AnnexB H264 on stdout; NAL units are split into frames on the fly.
* Video metadata is parsed from ffmpeg stderr during init.
*/
export async function demux(
input: string | PassThrough,
_opts: { format: string },
opts: { format: string; frameRate?: number },
): Promise<{
video: DemuxedStream | undefined;
audio: DemuxedStream | undefined;
close: () => void;
}> {
const _label = randomUUID();
const vPipe = new PassThrough({ objectMode: true, highWaterMark: 128 });
const aPipe = new PassThrough({ objectMode: true, highWaterMark: 128 });
// For stream input, spool to a temp file first so ffprobe can inspect it
// (ffprobe needs a seekable file; pipes can't be re-read). The stream is
// fully consumed before ffmpeg starts — acceptable for screen-share
// sources which are already fully buffered by yt-dlp in practice.
let spoolPath: string | null = null;
const cleanupSpool = () => {
if (spoolPath) {
import("node:fs").then(({ unlink }) => unlink(spoolPath!, () => {}));
spoolPath = null;
}
};
let effectiveInput: string;
if (typeof input === "string") {
effectiveInput = input;
} else {
spoolPath = join(tmpdir(), `golive-demux-${_label}.h264`);
const ws = createWriteStream(spoolPath);
await new Promise<void>((resolve, reject) => {
input.pipe(ws);
input.on("error", reject);
ws.on("finish", resolve);
ws.on("error", reject);
});
effectiveInput = spoolPath;
}
// Probe for codec + dimensions
let streams: Array<Record<string, unknown>> = [];
try {
streams = await probeStreams(effectiveInput);
} catch (_e) {
// probe failed (e.g. raw h264 without container) — infer h264 default
streams = [];
}
const v = streams.find((s) => s.codec_type === "video");
const a = streams.find((s) => s.codec_type === "audio");
let vInfo: DemuxedStream | undefined;
let aInfo: DemuxedStream | undefined;
if (v) {
const codecName = (v.codec_name as string) ?? "h264";
const rFrame = (v.r_frame_rate as string) ?? "0/1";
const [num, den] = rFrame.split("/").map((n) => Number(n));
vInfo = {
codec:
AVCodecID[
(codecName.toUpperCase() as keyof typeof AVCodecID) ??
"AV_CODEC_ID_H264"
] ?? AVCodecID.AV_CODEC_ID_H264,
codecName,
width: (v.width as number) ?? 0,
height: (v.height as number) ?? 0,
framerate_num: num ?? 0,
framerate_den: den ?? 1,
sample_rate: 0,
stream: vPipe,
};
} else {
// Probe failed (e.g. raw AnnexB h264 input) — still emit frames on the
// video pipe; playStream infers dimensions from the first frame.
vInfo = {
codec: AVCodecID.AV_CODEC_ID_H264,
codecName: "h264",
width: 0,
height: 0,
framerate_num: 0,
framerate_den: 1,
sample_rate: 0,
stream: vPipe,
};
}
if (a) {
const codecName = (a.codec_name as string) ?? "opus";
aInfo = {
codec:
AVCodecID[
(codecName.toUpperCase() as keyof typeof AVCodecID) ??
"AV_CODEC_ID_OPUS"
],
codecName,
width: 0,
height: 0,
framerate_num: 0,
framerate_den: 0,
sample_rate: Number(a.sample_rate) ?? 0,
stream: aPipe,
};
}
// Spawn ffmpeg — extract raw video (AnnexB for H264) to stdout
const isStream = typeof input !== "string";
const args: string[] = [
"-hide_banner",
// info level: stream init lines ("Stream #0:0: Video: h264...") go to
// stderr and are parsed for dimensions/fps.
"-loglevel",
"error",
"info",
// Input format hint: prepareStream always emits raw AnnexB H264 on
// pipe:0. Raw H264 has NO magic header, so ffmpeg's auto-detection
// fails with "Invalid data found when processing input" whenever the
// first bytes arrive late/buffered. Pin the demuxer input format.
...(isStream ? ["-f", "h264"] : []),
"-i",
effectiveInput,
isStream ? "pipe:0" : input,
"-c:v",
"copy",
"-an", // no audio in this minimal demuxer
@@ -262,25 +178,165 @@ export async function demux(
"h264",
"pipe:1",
];
const proc = spawn(FFMPEG, args, {
stdio: isStream ? ["pipe", "pipe", "pipe"] : ["ignore", "pipe", "pipe"],
});
console.log(
`[goLive:Demuxer] spawn ffmpeg pid=${proc.pid} input=${isStream ? "stream" : input} args=${args.join(" ")}`,
);
const proc = spawn(FFMPEG, args, { stdio: ["ignore", "pipe", "pipe"] });
// Pipe live input straight into ffmpeg stdin — never await stream end.
if (isStream && proc.stdin) {
input.pipe(proc.stdin);
input.on("error", () => proc.stdin?.destroy());
}
// Scan stdout for NAL units. Each NAL unit (between start codes) is one frame
// payload. We emit them individually; the packetizer chain handles FU-A.
// Set true once stderr metadata has been parsed (see handler below).
let parsedMeta = false;
// Parse stream metadata from ffmpeg stderr as it arrives (first chunk has
// the init lines). Fall back to H264 defaults if parsing fails.
let vInfo: DemuxedStream = {
codec: AVCodecID.AV_CODEC_ID_H264,
codecName: "h264",
width: 0,
height: 0,
framerate_num: 0,
framerate_den: 1,
sample_rate: 0,
stream: vPipe,
};
let aInfo: DemuxedStream | undefined;
let stderrBuf = "";
if (proc.stderr) {
proc.stderr.on("data", (d: Buffer) => {
const text = d.toString();
stderrBuf = (stderrBuf + text).slice(-16384);
// Surface actionable lines: ffmpeg errors + stream init lines
if (/error|invalid|no such|failed|cannot|not found|unable/i.test(text)) {
console.log(
`[goLive:Demuxer] ffmpeg stderr: ${text.trim().split("\n").slice(0, 4).join(" | ")}`,
);
}
if (parsedMeta) return;
const streamRe = /Stream #0:(\d+): (Video|Audio): ([^,]+)/g;
let m: RegExpExecArray | null;
const found: Array<{ kind: string; codecRaw: string }> = [];
// biome-ignore lint/suspicious/noAssignInExpressions: regex loop idiom
while ((m = streamRe.exec(stderrBuf)) !== null) {
found.push({ kind: m[2], codecRaw: m[3] });
}
const v = found.find((s) => s.kind === "Video");
const a = found.find((s) => s.kind === "Audio");
if (!v && !a) return;
parsedMeta = true;
if (v) {
const codecName = v.codecRaw.split(" ")[0].toLowerCase();
const dim = /(\d{2,5})x(\d{2,5})/.exec(stderrBuf);
const fps = /(\d+(?:\.\d+)?) fps/.exec(stderrBuf);
vInfo = {
codec:
AVCodecID[
(codecName.toUpperCase() as keyof typeof AVCodecID) ??
"AV_CODEC_ID_H264"
] ?? AVCodecID.AV_CODEC_ID_H264,
codecName,
width: dim ? Number(dim[1]) : 0,
height: dim ? Number(dim[2]) : 0,
framerate_num: fps ? Math.round(Number(fps[1]) * 1000) : 0,
framerate_den: fps ? 1000 : 1,
sample_rate: 0,
stream: vPipe,
};
}
if (a) {
const codecName = a.codecRaw.split(" ")[0].toLowerCase();
const sr = /(\d+) Hz/.exec(stderrBuf);
aInfo = {
codec:
AVCodecID[
(codecName.toUpperCase() as keyof typeof AVCodecID) ??
"AV_CODEC_ID_OPUS"
] ?? AVCodecID.AV_CODEC_ID_OPUS,
codecName,
width: 0,
height: 0,
framerate_num: 0,
framerate_den: 0,
sample_rate: sr ? Number(sr[1]) : 0,
stream: aPipe,
};
}
});
}
// Wait (briefly) for ffmpeg to print its stream init lines on stderr so
// vInfo carries real dimensions/fps. The lines arrive with the first chunk
// — a short timeout covers slow starts; callers fall back to sensible
// defaults when width/height are 0 anyway.
await Promise.race([
new Promise<void>((resolve) => {
const check = setInterval(() => {
if (parsedMeta) {
clearInterval(check);
resolve();
}
}, 25);
}),
new Promise<void>((resolve) => setTimeout(resolve, 1500)),
]);
// Scan stdout for AnnexB NAL units and group them into ACCESS UNITS
// (one picture). Discord's H264 decoder requires a complete access unit —
// parameter sets + slice — inside a single RTP frame. Emitting each NAL
// as its own frame (SPS/PPS/SEI separate from the slice) makes the decoder
// unable to produce ANY picture: production showed a black GoLive tile
// despite frames flowing (5892B slices + 4B PPS + 33B SPS as separate
// frames, each with a near-zero RTP timestamp delta). We therefore buffer
// NALs and flush one frame per slice, prepending the parameter sets that
// precede it, and timestamp it as ONE frame at the video frame rate.
let videoBuf = Buffer.alloc(0);
let frameCount = 0;
let pendingNals: Buffer[] = [];
let pendingHasSlice = false;
let pendingIsKey = false;
// Raw H264 streams carry no timing info — ffmpeg's h264 demuxer guesses
// 25fps on stderr. Prefer the caller's explicit frameRate (the encode
// setting); it drives both RTP timestamp advance and pacing.
const videoFps =
opts.frameRate ?? (vInfo.framerate_num / vInfo.framerate_den || 30);
const emitFrame = (nal: Uint8Array, isKeyFrame: boolean) => {
const flushAccessUnit = () => {
if (pendingNals.length === 0) return;
// AnnexB access unit: 00 00 00 01 + NAL for every buffered NAL. The
// packetizer (H264RtpPacketizer, StartSequence separator) needs the
// start codes to find NAL boundaries inside the frame.
const parts: Buffer[] = [];
for (const n of pendingNals) parts.push(startCode4, n);
const au = Buffer.concat(parts);
const isKey = pendingIsKey;
pendingNals = [];
pendingHasSlice = false;
pendingIsKey = false;
vPipe.write({
data: Buffer.from(nal),
data: au,
// One frame at videoFps: duration=1 in a 1/fps timebase →
// BaseMediaStream computes frametime=1000/fps ms → the RTP timestamp
// advances clockRate/fps per frame (3000 @ 30fps / 90kHz), which is
// what Discord's receiver expects for real-time video.
pts: frameCount,
duration: 1,
timeBase: { num: 1, den: 90000 },
flags: isKeyFrame ? AV_PKT_FLAG_KEY : 0,
timeBase: { num: 1, den: videoFps },
flags: isKey ? AV_PKT_FLAG_KEY : 0,
streamIndex: 0,
free: () => {},
});
frameCount++;
if (frameCount === 1 || frameCount % 30 === 0) {
console.log(
`[goLive:Demuxer] frames=${frameCount} last=${au.length}B key=${isKey}`,
);
}
};
if (proc.stdout) {
@@ -327,8 +383,21 @@ export async function demux(
while (end > 0 && nal[end - 1] === 0) end--;
if (end > 0) {
const nalTrimmed = nal.subarray(0, end);
const isIdr = (nalTrimmed[0] & 0x1f) === 5; // IDR
emitFrame(nalTrimmed, isIdr);
const nalType = nalTrimmed[0] & 0x1f;
const isSlice = nalType === 1 || nalType === 5;
if (isSlice) {
// A new slice while one is pending closes the previous
// access unit (x264 emits one slice per frame).
if (pendingHasSlice) flushAccessUnit();
pendingNals.push(Buffer.from(nalTrimmed));
pendingHasSlice = true;
if (nalType === 5) pendingIsKey = true;
} else {
// Parameter-set / SEI / AUD / filler NAL. After a slice these
// belong to the NEXT access unit — flush the completed frame.
if (pendingHasSlice) flushAccessUnit();
pendingNals.push(Buffer.from(nalTrimmed));
}
}
}
// Skip the 00 00 01 at scPos-3 to find next
@@ -346,21 +415,12 @@ export async function demux(
}
});
proc.stdout.on("end", () => {
if (videoBuf.length > 0) {
let end = videoBuf.length;
while (end > 0 && videoBuf[end - 1] === 0) end--;
if (end > 0) emitFrame(videoBuf.subarray(0, end), false);
}
flushAccessUnit();
vPipe.end();
aPipe.end();
});
}
if (proc.stderr) {
proc.stderr.on("data", () => {
/* errors swallowed */
});
}
proc.on("close", () => {
vPipe.end();
aPipe.end();
@@ -370,7 +430,6 @@ export async function demux(
proc.kill("SIGTERM");
vPipe.end();
aPipe.end();
cleanupSpool();
};
return { video: vInfo, audio: aInfo, close };
@@ -26,13 +26,24 @@ export function software(
} = {},
): () => EncoderSet {
const { x264, x265 } = opts;
const { preset: x264Preset = "superfast", tune: x264Tune = "film" } =
const { preset: x264Preset = "superfast", tune: x264Tune = "zerolatency" } =
x264 ?? {};
const { preset: x265Preset = "superfast", tune: x265Tune } = x265 ?? {};
return () => ({
H264: {
name: "libx264",
options: ["-forced-idr 1", `-tune ${x264Tune}`, `-preset ${x264Preset}`],
// -profile:v baseline is REQUIRED: the SDP advertises
// profile-level-id=42e01f (constrained baseline) and Discord's
// receiver decodes with that profile. x264's default is High — a
// High-profile bitstream against a baseline SDP negotiation fails to
// decode → black GoLive tile (production bug, fixed 2026-08-12).
// zerolatency matches @dank074 (no lookahead — correct for live).
options: [
"-forced-idr 1",
"-profile:v baseline",
`-tune ${x264Tune}`,
`-preset ${x264Preset}`,
],
},
H265: {
name: "libx265",
+102 -31
View File
@@ -48,7 +48,19 @@ export class Streamer {
this._client = client;
// listen for gateway dispatch events
this.client.on("raw", (packet) => {
this._gatewayEmitter.emit(packet.t, packet.d);
const t = packet.t as string;
if (
t === "STREAM_CREATE" ||
t === "STREAM_SERVER_UPDATE" ||
t === "VOICE_STATE_UPDATE" ||
t === "VOICE_SERVER_UPDATE"
) {
console.log(
`[goLive:Streamer] raw dispatch ${t}`,
JSON.stringify(packet.d).slice(0, 220),
);
}
this._gatewayEmitter.emit(t, packet.d);
});
}
@@ -65,6 +77,11 @@ export class Streamer {
}
sendOpcode(code: number, data: unknown): void {
// Direct instrumentation — bypasses the bootstrap debug filter (which
// drops messages without [VOICE / [ffmpeg / error / stream).
console.log(
`[goLive:Streamer] sendOpcode op=${code} d=${JSON.stringify(data)}`,
);
this.client.ws.broadcast({ op: code, d: data });
}
@@ -141,7 +158,6 @@ export class Streamer {
);
return;
}
this.signalStream();
const {
guildId: clientGuildId,
channelId: clientChannelId,
@@ -155,40 +171,85 @@ export class Streamer {
clientUserId,
clientChannelId,
(conn) => {
clearTimeout(streamTimeout);
clearInterval(retryInterval);
resolve(conn);
},
);
this.voiceConnection.streamConnection = streamConn;
this._gatewayEmitter.on(
"STREAM_CREATE",
(d: { stream_key: string; rtc_server_id: string }) => {
const { channelId, guildId, userId } = parseStreamKey(d.stream_key);
if (
clientGuildId !== guildId ||
clientChannelId !== channelId ||
clientUserId !== userId
) {
return;
}
streamConn.serverId = d.rtc_server_id;
streamConn.streamKey = d.stream_key;
streamConn.setSession(session_id);
},
);
this._gatewayEmitter.on(
"STREAM_SERVER_UPDATE",
(d: { stream_key: string; endpoint: string; token: string }) => {
const { channelId, guildId, userId } = parseStreamKey(d.stream_key);
if (
clientGuildId !== guildId ||
clientChannelId !== channelId ||
clientUserId !== userId
) {
return;
}
streamConn.setTokens(d.endpoint, d.token);
},
// Attach listeners BEFORE the first signal so a fast dispatch can't
// be lost between signalStream() and listener registration.
const onStreamCreate = (d: {
stream_key: string;
rtc_server_id: string;
}) => {
const { channelId, guildId, userId } = parseStreamKey(d.stream_key);
if (
clientGuildId !== guildId ||
clientChannelId !== channelId ||
clientUserId !== userId
) {
return;
}
streamConn.serverId = d.rtc_server_id;
streamConn.streamKey = d.stream_key;
streamConn.setSession(session_id);
};
const onStreamServerUpdate = (d: {
stream_key: string;
endpoint: string;
token: string;
}) => {
const { channelId, guildId, userId } = parseStreamKey(d.stream_key);
if (
clientGuildId !== guildId ||
clientChannelId !== channelId ||
clientUserId !== userId
) {
return;
}
streamConn.setTokens(d.endpoint, d.token);
};
this._gatewayEmitter.on("STREAM_CREATE", onStreamCreate);
this._gatewayEmitter.on("STREAM_SERVER_UPDATE", onStreamServerUpdate);
const cleanup = () => {
clearTimeout(streamTimeout);
clearInterval(retryInterval);
this._gatewayEmitter.removeListener("STREAM_CREATE", onStreamCreate);
this._gatewayEmitter.removeListener(
"STREAM_SERVER_UPDATE",
onStreamServerUpdate,
);
};
const streamTimeout = setTimeout(() => {
cleanup();
reject(
new Error(
"Timed out waiting for STREAM_CREATE/STREAM_SERVER_UPDATE from Discord (stream handshake) — voice media session may not be active",
),
);
}, 12_000);
// Discord sometimes drops the STREAM_CREATE request silently (upstream
// issue #217/#219) — resend a few times instead of giving up after one.
let attempt = 0;
const retryInterval = setInterval(() => {
attempt += 1;
if (attempt >= 4) {
clearInterval(retryInterval);
return;
}
console.log(
`[goLive:Streamer] createStream: retrying STREAM_CREATE (attempt ${attempt + 1}/4)`,
);
this.signalStream();
}, 3_000);
console.log(
`[goLive:Streamer] createStream: sending STREAM_CREATE (attempt 1/4)`,
);
this.signalStream();
});
}
@@ -241,6 +302,16 @@ export class Streamer {
channelId: channel_id,
botId: user_id,
} = this.voiceConnection;
// Un-deafen before requesting the stream (mimic real client). Do NOT
// set self_video: true — that flips on the bot's camera in Discord
// (visible to everyone); screen share should not enable the camera.
this.sendOpcode(GatewayOpCodes.VOICE_STATE_UPDATE, {
guild_id,
channel_id,
self_mute: false,
self_deaf: false,
self_video: false,
});
this.sendOpcode(GatewayOpCodes.STREAM_CREATE, {
type,
guild_id,
@@ -68,6 +68,7 @@ export class WebRtcConnWrapper {
private _audioTrack: NativeTrack | null = null;
private _videoTrack: NativeTrack | null = null;
private _videoCodec: WebRtcVideoCodec | null = null;
private _videoFrameLog = 0;
/** Assigned by BaseMediaConnection to send the gathered SDP to Discord. */
onLocalDescription: ((sdp: string) => void) | null = null;
@@ -114,7 +115,15 @@ export class WebRtcConnWrapper {
}
sendVideoFrame(frame: Buffer, frametime: number): void {
if (!this.ready || !this._videoTrack) return;
if (!this.ready || !this._videoTrack) {
if (this._videoFrameLog === 0) {
console.log(
`[goLive:WebRtc] sendVideoFrame DROPPED ready=${this.ready} track=${this._videoTrack !== null}`,
);
this._videoFrameLog++;
}
return;
}
const clockRate = CodecPayloadType[this._videoCodec ?? "H264"].clockRate;
if (this._videoCodec === "H264") {
let spsRewritten = false;
@@ -157,6 +166,12 @@ export class WebRtcConnWrapper {
}
this._videoTrack.sendFrame(frame);
this._videoTrack.addTimestamp(Math.round((frametime * clockRate) / 1000));
this._videoFrameLog++;
if (this._videoFrameLog === 1 || this._videoFrameLog % 30 === 0) {
console.log(
`[goLive:WebRtc] sendVideoFrame #${this._videoFrameLog} bytes=${frame.length} ready=${this.ready}`,
);
}
}
setPacketizer(videoCodec: string): void {
@@ -206,9 +206,19 @@ export function prepareStream(
: spawn(FFMPEG_BIN, args, { stdio: ["pipe", "pipe", "pipe"] });
if (proc.stdin && !isUrl) {
input.on("data", (chunk: Buffer) => proc.stdin?.write(chunk));
input.on("end", () => proc.stdin?.end());
input.on("error", () => proc.stdin?.destroy());
// Race guard: the merge ffmpeg may have already exited (transient 403
// or stream death) before this function attaches its listeners — the
// input's 'end'/'error' events then fire into the void and the encoder
// stdin NEVER receives EOF, leaving an encoder that waits forever and a
// screen share that shows a black tile with zero frames. Check the
// terminal state eagerly and EOF the encoder immediately.
if (input.readableEnded || input.destroyed) {
proc.stdin.end();
} else {
input.on("data", (chunk: Buffer) => proc.stdin?.write(chunk));
input.on("end", () => proc.stdin?.end());
input.on("error", () => proc.stdin?.destroy());
}
}
proc.stdout?.pipe(output);
@@ -262,15 +272,24 @@ export async function playStream(
options: PlayStreamOptions = {},
): Promise<void> {
const conn = await streamer.createStream();
console.log("[goLive:playStream] createStream resolved");
const { video, close: demuxClose } = await demux(prepared.output, {
format: options.format ?? "nut",
frameRate:
typeof options.frameRate === "number" ? options.frameRate : undefined,
});
console.log(
`[goLive:playStream] demux done codec=${video?.codecName ?? "?"} ${video?.width ?? 0}x${video?.height ?? 0} fps=${video ? video.framerate_num / video.framerate_den || 30 : 30}`,
);
if (!video) throw new Error("No video stream in media");
conn.setPacketizer(video.codecName);
conn.mediaConnection.setSpeaking(true);
console.log(
`[goLive:playStream] setPacketizer(${video.codecName}) + setSpeaking done`,
);
const w =
typeof options.width === "function"
@@ -310,14 +329,81 @@ export async function playStream(
}
};
return new Promise<void>((resolve) => {
vStream.once("finish", () => {
cleanup();
// First-frame watchdog: if the encoder never delivers a single frame
// (dead merge input, empty stream, codec mismatch), fail fast instead of
// "playing" a black tile forever. The demuxer resolves with fallback
// metadata even when no frame ever arrives, so this timeout is the only
// place that detects "started but nothing flowing".
let firstFrameTimer: NodeJS.Timeout | null = null;
let gotFirstFrame = false;
const firstFrame = new Promise<void>((resolve, reject) => {
firstFrameTimer = setTimeout(() => {
if (!gotFirstFrame) {
cleanup();
reject(
new Error(
"No video frames within 10s of stream start — input stream failed",
),
);
}
}, 10000);
video.stream.once("data", () => {
gotFirstFrame = true;
if (firstFrameTimer) clearTimeout(firstFrameTimer);
resolve();
});
});
return new Promise<void>((resolve, reject) => {
let settled = false;
const settle = (fn: () => void) => () => {
if (settled) return;
settled = true;
if (firstFrameTimer) clearTimeout(firstFrameTimer);
fn();
};
vStream.once("finish", () => {
settle(() => {
cleanup();
if (!gotFirstFrame) {
reject(new Error("Screen video stream ended without any frame"));
} else {
resolve();
}
})();
});
vStream.once("error", () => {
cleanup();
resolve();
settle(() => {
cleanup();
if (!gotFirstFrame) {
reject(new Error("Screen video stream errored before first frame"));
} else {
resolve();
}
})();
});
// The stream may end without ever producing a frame (input was
// silently dead) — surface that instead of resolving "successfully".
video.stream.once("end", () => {
settle(() => {
cleanup();
if (!gotFirstFrame) {
reject(new Error("Screen video stream ended before any frame"));
} else {
resolve();
}
})();
});
// Watchdog timeout: no frame arrived within 10s — fail fast instead of
// "playing" a black tile forever. cleanup() kills the encoder so the
// vStream finish/error handlers above still fire, but the settled guard
// ensures this rejection wins.
firstFrame.catch((err) => {
settle(() => {
cleanup();
reject(err);
})();
});
});
}
@@ -12,6 +12,7 @@ export const SYSTEM_RULES = `Kamu adalah asisten moderasi konten untuk server Di
## Normalisasi & Pertahanan Lintas Bahasa (WAJIB)
1. Campuran bahasa (Inggris/Indonesia/daerah) WAJIB dinormalisasi mental ke Bahasa Indonesia sebelum menilai intent. Jangan longgar hanya karena sintaksis campur (Polyglot Obfuscation).
2. Lakukan Named Entity Recognition agresif — nama orang/karakter (mis. "ren" setelah kata archaic "diagem") tetap dikenali sebagai nama.
3. <term_glossary> (bila ada) = definisi kata/slang/jargon yang tidak umum. Baca dulu arti kata yang tidak kamu kenal dari sana — jangan menebak dari bunyi/kemiripan. Kata yang tampak mencurigakan namun ternyata bermakna netral di glossary = AMAN; kata asing yang ternyata vulgar/terlarang di glossary = FLAG.
## Aturan Umum (AMAN — jangan flag)
- Slang: anjay, wkwk, gws, gaskeun, santuy, njir, baka, woy/woi, hadeh, astaga = AMAN.
@@ -73,6 +74,7 @@ RENDAH: harassment, vulgar_language terarah, offensive_username (Scunthorpe: "Sa
## Web Sebagai Bukti Utama
- <web_searches> ADALAH BUKTI UTAMA. Jika ada, WAJIB pakai hasilnya (hentai/scam/narkoba → flag; aman → clean). JANGAN abaikan. Jika tidak ada → gunakan pengetahuan internal.
- <term_glossary> = REFERENSI ARTI KATA, bukan bukti pelanggaran. Dipakai untuk memahami istilah yang tidak dikenal sebelum memutuskan.
- Prioritas bukti: <web_searches> > <web_content> > <media_analysis> > pengetahuan internal. <web_content> (URL fetch): gunakan isi, jangan flag hanya dari domain name.
## Pohon Keputusan
@@ -109,6 +109,7 @@ export function buildSystemPrompt(options: BuildSystemPromptOptions): string {
`- <conversation_context> = obrolan SEBELUM pesan target. Baris "[context]" di dalamnya BUKAN yang dinilai.\n` +
`- <user_profiles> = peta ringkasan kepribadian per user_id (attr as_of = kapan profil terakhir dibuat — profil lama mungkin tidak mencerminkan perilaku terkini); setiap <message> merujuk lewat <user_profile_ref user_id="..."/>.\n` +
`- <web_searches> / <web_content> = bukti web (lihat "Web Sebagai Bukti Utama").\n` +
`- <term_glossary> = kamus istilah: definisi kata/slang/jargon yang jarang dikenal (hasil pencarian Wikipedia via SearXNG). Gunakan untuk memahami arti kata yang tidak kamu kenal — JANGAN menebak atau mengarang arti.\n` +
`- <messages_to_analyze> = pesan-pesan TARGET yang WAJIB dinilai. Atribut <message>: id, user (nama server), time (ISO — kapan pesan dikirim), repetitions (N = teks pendek sama muncul N kali di batch — sinyal spam), bot (true jika dari bot), edited (true jika konten adalah hasil edit setelah posting).`,
);
@@ -12,6 +12,42 @@ const CACHE_PREFIX = "searxng:";
let redis: Redis | null = null;
/**
* Exposes the shared SearXNG Redis connection so other modules (e.g. the
* term glossary) reuse the same connection and cache prefix instead of
* opening their own. Returns null when Redis is unavailable.
*/
export function getSearxngRedis(): Redis | null {
return redis;
}
/** Builds a namespaced SearXNG cache key (shared across modules). */
export function makeSearxngCacheKey(namespace: string, key: string): string {
return `${CACHE_PREFIX}${namespace}:${key.toLowerCase().trim()}`;
}
/** Reads a value from the SearXNG Redis cache; null on miss/unavailable. */
export async function searxngCacheGet(key: string): Promise<string | null> {
if (!redis) return null;
try {
return await redis.get(key);
} catch {
return null;
}
}
/** Writes a value to the SearXNG Redis cache, fire-and-forget. */
export function searxngCacheSet(
key: string,
value: string,
ttlSeconds: number,
): void {
if (!redis) return;
redis.setex(key, ttlSeconds, value).catch(() => {
// Cache write failed silently
});
}
/**
* Initialize Redis connection for SearXNG cache.
* Safe to call multiple times — only creates one connection.
@@ -51,19 +87,26 @@ export interface SearxngResult {
/**
* Search SearXNG for a query and return structured results.
* Uses Redis cache when available — same query within 24h returns cached results.
*
* @param engines Optional comma-separated SearXNG engine list to constrain
* the search (e.g. "wikipedia"). When set, results are cached under a
* separate cache namespace so engine-specific results never collide.
*/
export async function searchSearxng(
query: string,
category: "general" | "news" | "science" = "general",
engines?: string,
timeoutMs: number = TIMEOUT_MS,
): Promise<SearxngResult[]> {
const cacheKey = `${CACHE_PREFIX}${category}:${query.toLowerCase().trim()}`;
const engineNs = engines ? `eng:${engines}` : "auto";
const cacheKey = makeSearxngCacheKey(`${category}:${engineNs}`, query);
// Try cache first
if (redis) {
try {
const cached = await redis.get(cacheKey);
if (cached) {
log.debug({ query, category }, "SearXNG cache HIT");
log.debug({ query, category, engines }, "SearXNG cache HIT");
return JSON.parse(cached) as SearxngResult[];
}
} catch {
@@ -73,8 +116,11 @@ export async function searchSearxng(
// Cache miss — hit SearXNG API
try {
const url = `${SEARXNG_BASE_URL}/search?q=${encodeURIComponent(query)}&format=json&language=id&categories=${category}`;
const { controller, clear } = createAbortControllerWithTimeout(TIMEOUT_MS);
const engineParam = engines
? `&engines=${encodeURIComponent(engines)}`
: "";
const url = `${SEARXNG_BASE_URL}/search?q=${encodeURIComponent(query)}&format=json&language=id&categories=${category}${engineParam}`;
const { controller, clear } = createAbortControllerWithTimeout(timeoutMs);
try {
const response = await fetch(url, {
@@ -0,0 +1,419 @@
/**
* termGlossary.ts
*
* Per-word "kamus" enrichment for LLM moderation.
*
* Problem: the moderation LLM often meets words it does not know — regional
* slang (Jawa/Sunda), foreign terms, niche anime/game jargon, or obscure
* technical vocabulary. When it guesses, it either invents a wrong meaning
* (false positive on a safe word) or misses a violation hidden in unfamiliar
* wording (false negative on an unknown vulgar/slang term).
*
* Solution: extract candidate "unknown-looking" words from message content,
* look each one up on Wikipedia via SearXNG, and inject the definitions into
* the LLM prompt as a `<term_glossary>` block so verdicts are based on facts
* instead of guesses.
*
* Cost control:
* - definitions are cached in an in-memory LRU AND in Redis (shared with the
* SearXNG cache) — a term is searched at most once per TTL across the whole
* service, so repeat lookups are effectively free;
* - lookups per batch are bounded (AI_GLOSSARY_MAX_TERMS);
* - Wikipedia-only search first, generic search as a fallback;
* - everything degrades gracefully: no Redis, no SearXNG, no Wikipedia match
* → the block is simply omitted and moderation proceeds as before.
*/
import { LRUCache } from "lru-cache";
import { createChildLogger } from "@/shared/logger/index";
import { config } from "../../shared/config/config.js";
import { escapeXml } from "./moderationBuilders.js";
import {
makeSearxngCacheKey,
searchSearxng,
searxngCacheGet,
searxngCacheSet,
} from "./searxngSearch.js";
const log = createChildLogger("term-glossary");
// ---------------------------------------------------------------------------
// Constants
// ---------------------------------------------------------------------------
/** Redis TTL for a successfully resolved definition (definitions are stable). */
const DEF_TTL_SECONDS = 7 * 24 * 60 * 60;
/** Redis TTL for a lookup that found nothing — don't re-search every batch. */
const MISS_TTL_SECONDS = 24 * 60 * 60;
/** Sentinel stored in caches for "term has no resolvable definition". */
const EMPTY_SENTINEL = "__not_found__";
/** Per-search timeout — keep glossary lookups snappy even on a slow SearXNG. */
const GLOSSARY_SEARCH_TIMEOUT_MS = 5000;
/** Max definition snippet length kept in the prompt. */
const MAX_DEFINITION_CHARS = 300;
/** In-memory cache: term (lowercase) → definition | NOT_FOUND sentinel. */
const NOT_FOUND: TermDefinition = {
term: "__not_found__",
definition: "",
sourceUrl: "",
};
const termLru = new LRUCache<string, TermDefinition>({
max: 2000,
ttl: 24 * 60 * 60 * 1000,
});
// ---------------------------------------------------------------------------
// Term extraction
// ---------------------------------------------------------------------------
/** Word tokenizer — letters/digits plus internal -_'· (handles "well-known",
* "node_modules", diacritics). */
const WORD_RE = /[\p{L}\p{N}]+(?:[-_'·][\p{L}\p{N}]+)*/gu;
/** Removes URLs, Discord mentions/custom emoji, code fences, markdown noise. */
function cleanContent(raw: string): string {
return raw
.replace(/https?:\/\/\S+/gi, " ")
.replace(/<@!?\d+>/g, " ")
.replace(/<#\d+>/g, " ")
.replace(/<a?:\w+:\d+>/g, " ")
.replace(/[`*_~|>\[\]]/g, " ")
.replace(/[\p{Emoji}\p{Extended_Pictographic}]/gu, " ")
.replace(/\s+/g, " ")
.trim();
}
/** Filters out tokens that are useless as glossary candidates (numbers,
* repeated-char noise, mega-tokens). */
function isNoiseWord(word: string): boolean {
if (word.length > 28) return true;
if (/^\d+$/.test(word)) return true;
const lower = word.toLowerCase();
// "aaaa…", "wwwwww" — single repeated character
if (/^(.)\1{2,}$/.test(lower)) return true;
// "wkwk", "hehe", "69" alternations — repeated 23 char base. "meme" is
// the one legit 4-letter word this matches; it is whitelisted below.
if (/^([a-z]{2,3})\1{1,}$/.test(lower)) return true;
return false;
}
/** Deterministic bonus for words that look like proper nouns or foreign. */
function scoreWord(word: string): number {
let score = 1;
// Capitalized first letter (proper noun / title) but not ALL-CAPS acronyms
if (/^[A-Z]/.test(word) && !/^[A-Z]{2,}$/.test(word)) score += 3;
// Contains a letter outside basic latin → regional/foreign spelling
if (/[\p{L}]/u.test(word.replace(/[A-Za-z]/g, ""))) score += 2;
// Contains an internal apostrophe or hyphen → likely a named entity
if (/[-_'’·]/.test(word)) score += 2;
return score;
}
const STOPWORDS = new Set(
// ── Bahasa Indonesia ────────────────────────────────────────────────
(
" yang dan di ke dari ini itu dengan untuk pada dalam adalah akan telah sudah bisa dapat harus tidak juga saya kamu kita kami mereka dia aku kau gua lu lo gw gue elu anda kalian nya kah lah pun ya yah kan sih dong deh kok loh toh aja saja gitu gini begitu begini tapi tetapi namun atau karena sebab jika kalau bila maka supaya agar meski meskipun walau walaupun ketika saat setelah sebelum selama antara terhadap tentang mengenai bagi oleh secara sebagai seperti daripada tanpa hingga sampai sejak menuju bahwa padahal sebenarnya sepertinya mungkin memang jadi lalu terus akhirnya misalnya contohnya banyak sedikit semua seluruh setiap tiap beberapa ada bukan jangan boleh mau ingin pengen nggak ngak gak ga kagak ngga ndak nanti kemarin besok hari ini sekarang waktu itu masih sedang belum pernah sering selalu kadang jarang cepat lambat awal akhir baru lama besar kecil tinggi rendah panjang pendek baik buruk benar salah sama beda penting biasanya selamat terima kasih makasih sangat sekali paling cuma cuman hanya lebih kurang sekitar hampir ternyata rupanya begitu gimana bagaimana kenapa mengapa siapa apa mana kapan darimana kemana bilang ngomong omong kata tadi dulu terus lagi tetap pasti seharusnya sebaiknya seakan seolah kayaknya keliatan kelihatan ketahuan disini disitu disana kesini kesana bener pake pakai kayak emang lagian mulu istilah istilahnya banget" +
// ── English ───────────────────────────────────────────────────────
" the a an and or but if then else for to in on at by with without from of is are was were be been being have has had do does did will would can could should may might must shall this that these those it its i you he she we they them their there here when where why how what which who whom whose only very just about above after before below under over into onto within upon against between among during through across along around behind beyond near off out up down now then so as not no yes ok okay" +
// ── Common net slang / acronyms the LLM already knows ──────────────
" lol omg wtf idk btw tbh imo aka fyi nsfw smh nvm asap afk brb gg wp ty np mb sry thx kk oke okk ygy frfr"
).split(/\s+/),
);
/**
* Words that are either already defined by the moderation rules, or are so
* common (brands, tech vocabulary, project names) that a Wikipedia lookup is
* a guaranteed miss/waste. Keeps the glossary focused on genuinely unknown
* terms.
*/
const KNOWN_SAFE_TERMS = new Set(
(
"discord youtube google facebook instagram twitter tiktok whatsapp telegram netflix spotify steam github gitlab bitbucket chatgpt openai anthropic claude deepseek gemini llama copilot cursor vscode vscodium jetbrains intellij pycharm webstorm sublime codeblocks" +
" docker kubernetes k8s linux ubuntu debian arch fedora manjaro kali windows macos android ios chrome firefox safari edge opera brave" +
" react nextjs next vue svelte angular node nodejs deno bun pnpm yarn npm javascript typescript python golang go rust java kotlin swift cplusplus cpp css html json xml yaml toml regex backend frontend database mysql postgres postgresql mongodb redis qdrant sqlite nosql graphql rest websocket webhook" +
" bug crash error debug fix issue pr merge commit push pull branch main master dev staging production server client app website web browser" +
" stream streaming video audio voice call camera screen share screenshare gameplay gaming game play steam epic xbox playstation nintendo switch console" +
" bot discordbot moderation moderator admin member user profile avatar channel server guild message chat dm reply forward embed sticker emoji role permission" +
" meme code coding ngoding programmer program developer engineer software hardware cpu gpu ram rom storage disk network internet wifi lan ip dns vpn proxy cloud aws azure gcp vercel netlify heroku railway render vps hosting domain ssl login logout register account password email username" +
" anime manga waifu husbando tsundere moe otaku wibu weeb otome isekai shonen seinen josei manga manhwa manhua doujin" +
" anjay wkwk wkwkwk gws gaskeun santuy njir baka woy woi hadeh astaga asu anjing bangsat ngehe asal alay lebay caper mabar" +
" asus bete imphnen impnhen ngab" +
" syahadat sholat shalat solat puasa zakat haji umrah doa tuhan nabi allah yesus muhammad hashem" +
" loli shota incest exhibition furry fursuit cosplay costume" +
" gaza palestine israel yahudi yahud israel palestina israeli" +
" hokkian mandarin arabic jawa sunda betawi minang bugis batak melayu inggris indonesia"
).split(/\s+/),
);
function isKnownTerm(word: string): boolean {
return STOPWORDS.has(word) || KNOWN_SAFE_TERMS.has(word);
}
/** True when a quoted phrase is mostly filler words (skip it). */
function isMostlyStopwords(phrase: string): boolean {
const words = phrase
.toLowerCase()
.split(/[^a-zà-öø-ÿ]+/i)
.filter(Boolean);
if (words.length === 0) return true;
const stopCount = words.filter((w) => STOPWORDS.has(w)).length;
return stopCount / words.length >= 0.6;
}
export interface ExtractGlossaryOptions {
maxTerms?: number;
minWordLength?: number;
}
/**
* Extracts candidate terms that the LLM might not know from message content.
* Returns at most `maxTerms` terms (default from config), scored by how
* "unknown-looking" they are (proper nouns, foreign spelling, quoted phrases).
*/
export function extractGlossaryTerms(
contents: string[],
options: ExtractGlossaryOptions = {},
): string[] {
const maxTerms = options.maxTerms ?? config.AI_GLOSSARY_MAX_TERMS;
const minWordLength =
options.minWordLength ?? config.AI_GLOSSARY_MIN_WORD_LENGTH;
const candidates = new Map<string, { word: string; score: number }>();
const push = (rawWord: string, score: number): void => {
const clean = rawWord
.trim()
.replace(/^[^\p{L}\p{N}]+|[^\p{L}\p{N}]+$/gu, "");
if (clean.length < minWordLength) return;
const key = clean.toLowerCase();
if (isKnownTerm(key) || isNoiseWord(clean)) return;
const existing = candidates.get(key);
if (existing) {
existing.score += score + 1;
} else {
candidates.set(key, { word: clean, score });
}
};
for (const content of contents) {
if (!content) continue;
const cleaned = cleanContent(content);
if (!cleaned) continue;
// Quoted phrases — explicit terms the user called out
for (const m of cleaned.matchAll(/"([^"]{2,80})"/g)) {
const phrase = m[1].trim();
const wordCount = phrase.split(/\s+/).length;
if (wordCount >= 2 && wordCount <= 6 && !isMostlyStopwords(phrase)) {
push(phrase, 10);
}
}
// Individual words
for (const m of cleaned.matchAll(WORD_RE)) {
const w = m[0];
if (w.length < minWordLength) continue;
if (isNoiseWord(w)) continue;
const key = w.toLowerCase();
if (isKnownTerm(key)) continue;
push(w, scoreWord(w));
}
}
return Array.from(candidates.values())
.sort((a, b) => b.score - a.score)
.slice(0, maxTerms)
.map((c) => c.word);
}
// ---------------------------------------------------------------------------
// Definition lookup (cached: LRU → Redis → SearXNG/Wikipedia)
// ---------------------------------------------------------------------------
export interface TermDefinition {
term: string;
definition: string;
sourceUrl: string;
}
/** Picks the best definition from search results, preferring Wikipedia. */
function pickDefinition(
results: Array<{ title: string; url: string; snippet: string }>,
term: string,
): TermDefinition | null {
const best = results.find((r) => /wikipedia/i.test(r.url)) ?? results[0];
if (!best) return null;
const snippet = (best.snippet || best.title || "").trim();
if (snippet.length < 10) return null;
const definition =
snippet.length > MAX_DEFINITION_CHARS
? `${snippet.slice(0, MAX_DEFINITION_CHARS - 1).trimEnd()}`
: snippet;
return { term, definition, sourceUrl: best.url };
}
async function lookupTermDefinition(
term: string,
): Promise<TermDefinition | null> {
const key = term.toLowerCase().trim();
// 1. In-memory LRU — same process, instant
const lruHit = termLru.get(key);
if (lruHit) return lruHit === NOT_FOUND ? null : lruHit;
// 2. Redis — shared across processes/workers
const cacheKey = makeSearxngCacheKey("def", key);
const cached = await searxngCacheGet(cacheKey);
if (cached !== null) {
if (cached === EMPTY_SENTINEL) {
termLru.set(key, NOT_FOUND);
return null;
}
try {
const parsed = JSON.parse(cached) as {
definition?: string;
sourceUrl?: string;
};
if (parsed.definition) {
const def: TermDefinition = {
term,
definition: parsed.definition,
sourceUrl: parsed.sourceUrl ?? "",
};
termLru.set(key, def);
return def;
}
} catch {
// malformed cache entry — fall through to search
}
}
// 3. Live search — Wikipedia first, generic search as fallback
try {
let def = pickDefinition(
await searchSearxng(
key,
"general",
"wikipedia",
GLOSSARY_SEARCH_TIMEOUT_MS,
),
term,
);
if (!def) {
def = pickDefinition(
await searchSearxng(
`${key} definisi arti`,
"general",
undefined,
GLOSSARY_SEARCH_TIMEOUT_MS,
),
term,
);
}
if (def) {
searxngCacheSet(
cacheKey,
JSON.stringify({
definition: def.definition,
sourceUrl: def.sourceUrl,
}),
DEF_TTL_SECONDS,
);
termLru.set(key, def);
log.debug({ term: key }, "Term glossary resolved definition");
return def;
}
} catch (err) {
log.debug(
{ term: key, error: err instanceof Error ? err.message : String(err) },
"Term glossary lookup failed — skipping term",
);
}
// No definition — cache the miss so we do not re-search every batch.
searxngCacheSet(cacheKey, EMPTY_SENTINEL, MISS_TTL_SECONDS);
termLru.set(key, NOT_FOUND);
return null;
}
/**
* Looks up definitions for a batch of terms, in parallel. Returns a map of
* term → definition for the terms that resolved. Errors/misses are skipped.
*/
export async function lookupTermDefinitions(
terms: string[],
): Promise<Map<string, TermDefinition>> {
const map = new Map<string, TermDefinition>();
if (terms.length === 0) return map;
const results = await Promise.allSettled(terms.map(lookupTermDefinition));
for (let i = 0; i < terms.length; i++) {
const r = results[i];
if (r.status === "fulfilled" && r.value) {
map.set(r.value.term, r.value);
}
}
return map;
}
// ---------------------------------------------------------------------------
// Prompt formatting
// ---------------------------------------------------------------------------
/**
* Formats definitions as a `<term_glossary>` XML block for the LLM prompt:
*
* <term_glossary>
* <term word="ngab" source="https://…">definisi…</term>
* </term_glossary>
*
* Returns "" when there are no definitions (the block is then omitted).
*/
export function formatTermGlossary(
defs: ReadonlyMap<string, TermDefinition>,
): string {
if (!defs || defs.size === 0) return "";
const lines = Array.from(defs.values()).map(
(d) =>
` <term word="${escapeXml(d.term)}" source="${escapeXml(d.sourceUrl)}">${escapeXml(d.definition)}</term>`,
);
return `<term_glossary>\n${lines.join("\n")}\n</term_glossary>`;
}
// ---------------------------------------------------------------------------
// Convenience: full pipeline
// ---------------------------------------------------------------------------
export interface GlossaryBlockOptions extends ExtractGlossaryOptions {
enabled?: boolean;
}
/**
* One-shot helper: extract terms from message contents, look up definitions,
* and return the formatted `<term_glossary>` block ("" when disabled or no
* definitions found). Safe to call on every batch — cached lookups make it
* cheap.
*/
export async function buildTermGlossaryBlock(
contents: string[],
options: GlossaryBlockOptions = {},
): Promise<string> {
const enabled = options.enabled ?? config.AI_GLOSSARY_ENABLED;
if (!enabled) return "";
if (contents.length === 0) return "";
const terms = extractGlossaryTerms(contents, options);
if (terms.length === 0) return "";
const defs = await lookupTermDefinitions(terms);
if (defs.size === 0) return "";
const block = formatTermGlossary(defs);
log.debug(
{ terms: terms.length, definitions: defs.size },
"Term glossary block built",
);
return block;
}
@@ -37,6 +37,7 @@ import {
formatSearchResults,
searchSearxng,
} from "./searxngSearch.js";
import { buildTermGlossaryBlock } from "./termGlossary.js";
import { getRecentCorrectedModerations } from "./textCacheStore.js";
import { extractUrlsFromText, fetchUrlSafely } from "./urlFetcher.js";
import { getUserProfile } from "./userProfileStore.js";
@@ -145,9 +146,17 @@ export async function runTextOnlyBatch(
return map;
})();
const [urlFetchMaps, searxngResults] = await Promise.all([
// Term glossary — per-word Wikipedia lookups for words the LLM may not
// know (slang, jargon, regional language). Cached in Redis + in-memory, so
// repeat terms resolve instantly and only genuinely new words hit SearXNG.
const glossaryPromise = buildTermGlossaryBlock(
targets.map((msg) => getAnalysisContent(msg)),
).catch(() => "");
const [urlFetchMaps, searxngResults, glossaryBlock] = await Promise.all([
urlFetchPromise,
searxngPromise,
glossaryPromise,
]);
const urlFetchMap = urlFetchMaps.text;
@@ -368,6 +377,7 @@ export async function runTextOnlyBatch(
userProfilesBlock?.trimEnd() ?? "",
contextBlock?.trimEnd() ?? "",
searxngBlock,
glossaryBlock,
`<messages_to_analyze>\n${messagesBlock}\n</messages_to_analyze>`,
].filter((b) => b.trim().length > 0);
return {
@@ -87,6 +87,7 @@ import {
formatSearchResults,
searchSearxng,
} from "./searxngSearch.js";
import { buildTermGlossaryBlock } from "./termGlossary.js";
import { extractUrlsFromText } from "./urlFetcher.js";
import { getUserProfile } from "./userProfileStore.js";
import {
@@ -413,6 +414,12 @@ export async function prepareMediaMessage(
searxngXml = `\n<web_searches>\n${parts.join("\n")}\n</web_searches>`;
}
// Term glossary — cached per-word Wikipedia definitions for words the LLM
// may not know. Bounded and cached (in-memory + Redis), so this adds no
// meaningful latency to the media path either.
const glossaryXml = await buildTermGlossaryBlock([content]).catch(() => "");
const glossaryCtx = glossaryXml ? `\n${glossaryXml}` : "";
// Build XML block
const webTexts = webTextMap.get(targetId) ?? [];
const mediaAnalyses = mediaAnalysisMap.get(targetId) ?? [];
@@ -466,6 +473,6 @@ export async function prepareMediaMessage(
const isBot = resolveIsBot(target);
const isEdited = resolveIsEdited(target);
const messageBlock = `<message id="${escapeXml(target.id)}" user="${escapeXml(resolveDisplayName(target))}" time="${new Date(target.created_at).toISOString()}"${isBot ? ` bot="true"` : ""}${isEdited ? ` edited="true"` : ""}>\n ${repXml}${profileRef ? `\n ${profileRef}` : ""}${refXml ? `\n ${refXml}` : ""}\n <content>${escapeXml(truncateForAi(content))}</content>${mediaContext ? ` ${escapeXml(mediaContext)}` : ""}${webContext}${mediaAnalysisContext}${searxngXml}\n</message>`;
const messageBlock = `<message id="${escapeXml(target.id)}" user="${escapeXml(resolveDisplayName(target))}" time="${new Date(target.created_at).toISOString()}"${isBot ? ` bot="true"` : ""}${isEdited ? ` edited="true"` : ""}>\n ${repXml}${profileRef ? `\n ${profileRef}` : ""}${refXml ? `\n ${refXml}` : ""}\n <content>${escapeXml(truncateForAi(content))}</content>${mediaContext ? ` ${escapeXml(mediaContext)}` : ""}${webContext}${mediaAnalysisContext}${searxngXml}${glossaryCtx}\n</message>`;
return { targetId, messageBlock };
}
@@ -1,4 +1,7 @@
import { type ChildProcess, spawn } from "node:child_process";
import { chmodSync, mkdtempSync, rmSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { PassThrough, type Readable } from "node:stream";
import { StreamType } from "@discordjs/voice";
import { createChildLogger } from "@/shared/logger/index";
@@ -79,6 +82,22 @@ export function transcodeToHighQualityOgg(
);
input.pipe(proc.stdin);
// ffmpeg teardown closes stdin while the upstream source may still write —
// swallow EPIPE / destroyed-stream errors so they don't crash the gateway.
proc.stdin.on("error", (err: NodeJS.ErrnoException) => {
if (
err.code === "EPIPE" ||
err.code === "ERR_STREAM_DESTROYED" ||
err.code === "ERR_STREAM_WRITE_AFTER_END"
) {
logger.debug(
{ code: err.code },
"Transcode stdin closed during teardown",
);
} else {
logger.error({ error: err.message }, "Transcode stdin error");
}
});
activeProcesses.add(proc);
const cleanup = () => {
@@ -248,6 +267,20 @@ export function resolveMediaUrl(
// `--print` headers to stderr — pipe stdout immediately so the child
// never blocks on a full pipe while we wait for the headers on stderr.
const mediaStream = new PassThrough();
// Teardown (player stop / ffmpeg exit) destroys this stream while
// yt-dlp may still push bytes — without a listener an EPIPE /
// ERR_STREAM_DESTROYED surfaces as an uncaughtException.
mediaStream.on("error", (err: NodeJS.ErrnoException) => {
if (
err.code === "EPIPE" ||
err.code === "ERR_STREAM_DESTROYED" ||
err.code === "ERR_STREAM_WRITE_AFTER_END"
) {
logger.debug({ code: err.code }, "Media stream closed during teardown");
} else {
logger.error({ error: err.message }, "Media stream error");
}
});
proc.stdout.pipe(mediaStream);
let stderrBuf = "";
@@ -348,241 +381,97 @@ export function resolveMediaUrl(
* Resolve a media URL to a single playable input stream for screen share /
* GoLive streaming.
*
* yt-dlp `--get-url` with `bestvideo+bestaudio` prints the video-only and
* audio-only URLs on SEPARATE lines. The old code took only the first line
* (video-only) ffmpeg had no audio track GoLive stream had no sound.
* Streams the merged video+audio media directly from yt-dlp stdout (`-o -`).
*
* This returns a single input that `prepareStream` (which accepts only ONE
* ffmpeg input) can consume while STILL including audio:
* - If yt-dlp offers a merged progressive URL (one URL, video+audio) it is
* returned directly.
* - Otherwise the video-only + audio-only DASH URLs are fetched in the SAME
* yt-dlp run (signature URLs expire quickly) and merged locally by an
* ffmpeg process into a single NUT stream, which is streamed to the
* consumer over a Readable. NUT over stdin auto-probes cleanly (verified:
* av1+opus merge H264+opus transcode).
* This is deliberately NOT the old --dump-single-json + manual URL-fetch
* approach: YouTube signs DASH URLs for the extracting client and rejects
* them with 403 when fetched raw by ffmpeg/curl (verified 2026-08-12: even
* curl with the EXACT http_headers from the yt-dlp dump got 403 on some
* videos, while yt-dlp's own downloader succeeded). Streaming from yt-dlp
* lets it handle auth, cookies and transient retries internally the same
* mechanism resolveMediaUrl already uses for music playback.
*
* @returns a direct video URL (string) or a Readable of the merged NUT stream.
* @returns a Readable of the merged media stream.
*/
export function getDirectScreenInput(url: string): Promise<string | Readable> {
return new Promise<string | Readable>((resolve, reject) => {
export function getDirectScreenInput(url: string): Promise<Readable> {
return new Promise<Readable>((resolve) => {
// Merge fragments must NOT be written to the process CWD — the Nix
// store dir is read-only for the deployed gateway (EACCES). Use a
// per-run temp dir (world-writable like /tmp) so parallel/retry runs
// never collide on merge fragments and any user can write to it.
const tmpDir = mkdtempSync(join(tmpdir(), "gmw-ytdlp-"));
chmodSync(tmpDir, 0o1777);
const args = [
url,
"--dump-single-json",
"--format",
"-f",
"bestvideo[protocol^=http]+bestaudio[protocol^=http]/best[protocol^=http]/best",
"-o",
"-",
"--no-playlist",
"--no-warnings",
"--quiet",
// NOTE: deliberately NOT --no-simulate. Simulate mode still resolves the
// requested format URLs into the JSON (requested_formats[].url), and it
// avoids yt-dlp writing .part files into the process CWD — which is the
// read-only Nix store dir for the deployed gateway (EACCES).
"--no-progress",
"-P",
tmpDir,
url,
];
logger.info({ url }, "Spawning yt-dlp for screen share input resolution");
logger.info({ url }, "Spawning yt-dlp for screen share input streaming");
const proc = spawn("yt-dlp", args, {
stdio: ["pipe", "pipe", "pipe"],
stdio: ["ignore", "pipe", "pipe"],
});
activeProcesses.add(proc);
let stdoutBuf = "";
const stream = new PassThrough();
proc.stdout.pipe(stream);
let stderrBuf = "";
const MAX_STDERR = 4096;
const MAX_STDOUT = 8 * 1024 * 1024; // JSON metadata + requested format URLs
proc.stderr?.on("data", (chunk: Buffer) => {
if (stderrBuf.length < MAX_STDERR) {
stderrBuf += chunk.toString("utf8");
}
});
if (proc.stdout) {
proc.stdout.on("data", (chunk: Buffer) => {
if (stdoutBuf.length < MAX_STDOUT) {
stdoutBuf += chunk
.toString("utf8")
.slice(0, MAX_STDOUT - stdoutBuf.length);
}
});
}
if (proc.stderr) {
proc.stderr.on("data", (chunk: Buffer) => {
if (stderrBuf.length < MAX_STDERR) {
stderrBuf += chunk
.toString("utf8")
.slice(0, MAX_STDERR - stderrBuf.length);
}
});
}
let producedData = false;
stream.once("data", () => {
producedData = true;
});
proc.on("error", (err: NodeJS.ErrnoException) => {
activeProcesses.delete(proc);
rmSync(tmpDir, { recursive: true, force: true });
if (err.code === "ENOENT") {
reject(buildNotInstalledError());
stream.destroy(buildNotInstalledError());
} else {
reject(new Error(`yt-dlp failed to start: ${err.message}`));
stream.destroy(new Error(`yt-dlp failed to start: ${err.message}`));
}
});
proc.on("close", (code) => {
activeProcesses.delete(proc);
if (code !== 0) {
rmSync(tmpDir, { recursive: true, force: true });
// Fail fast: a download that dies before producing ANY bytes (e.g.
// transient YouTube 403) cannot feed the encoder — destroy the stream
// so the caller retries with a fresh yt-dlp run instead of streaming
// a silent black tile.
if (code !== 0 && !producedData && !stream.destroyed) {
const detail = stderrBuf.trim() ? `: ${stderrBuf.trim()}` : "";
reject(
stream.destroy(
new Error(
`yt-dlp screen input resolution exited with code ${code}${detail}`,
`yt-dlp screen input stream failed (exit ${code})${detail}`,
),
);
return;
}
let parsed: Record<string, unknown>;
try {
parsed = JSON.parse(stdoutBuf.trim()) as Record<string, unknown>;
} catch (parseErr) {
reject(
new Error(
`Failed to parse yt-dlp JSON for screen input: ${(parseErr as Error).message}`,
),
);
return;
}
resolveScreenInput(parsed).then(resolve, (err: unknown) => {
const message = err instanceof Error ? err.message : String(err);
reject(
new Error(`Failed to build screen input for "${url}": ${message}`),
);
});
});
// Resolve immediately — data flows as yt-dlp downloads. The caller's
// resolveInputWithRetry validates the first byte and retries on failure.
resolve(stream);
});
}
/**
* From a parsed yt-dlp JSON info dict, decide how to feed a single ffmpeg
* input with both video and audio.
*/
async function resolveScreenInput(
info: Record<string, unknown>,
): Promise<string | Readable> {
const requested = info.requested_formats as
| Array<Record<string, unknown>>
| undefined;
// Merged/progressive single URL (video+audio in one). Common when yt-dlp
// selects a single format (e.g. format 18 progressive mp4) or when a direct
// muxed URL is available.
const singleUrl = info.url as string | undefined;
const singleHasAudio =
info.acodec !== "none" &&
typeof info.acodec === "string" &&
info.acodec.length > 0;
if (typeof singleUrl === "string" && singleUrl && singleHasAudio) {
logger.debug("Screen share uses merged progressive single URL");
return singleUrl;
}
// Separate video-only + audio-only DASH formats → merge locally via ffmpeg.
if (Array.isArray(requested) && requested.length >= 2) {
const video = requested.find(
(rf) => rf.vcodec && String(rf.vcodec) !== "none",
);
const audio = requested.find(
(rf) => rf.acodec && String(rf.acodec) !== "none",
);
const videoUrl = video?.url as string | undefined;
const audioUrl = audio?.url as string | undefined;
if (
typeof videoUrl === "string" &&
videoUrl.length > 0 &&
typeof audioUrl === "string" &&
audioUrl.length > 0
) {
return mergeScreenStreams(videoUrl, audioUrl);
}
}
throw new Error(
"yt-dlp returned neither a merged progressive URL nor a video+audio format pair",
);
}
/**
* Merge a video-only URL and an audio-only URL into a single NUT stream using
* a child ffmpeg process. Both URLs come from the same yt-dlp run, so they
* share the same signature/expiry and are consumed immediately.
*/
function mergeScreenStreams(videoUrl: string, audioUrl: string): Readable {
logger.info("Merging video+audio DASH streams into a single NUT input");
const ffmpeg = spawn(
"ffmpeg",
[
"-hide_banner",
"-loglevel",
"error",
"-reconnect",
"1",
"-reconnect_streamed",
"1",
"-reconnect_delay_max",
"5",
"-i",
videoUrl,
"-i",
audioUrl,
"-map",
"0:v:0",
"-map",
"1:a:0",
"-c:v",
"copy",
"-c:a",
"copy",
"-f",
"nut",
"pipe:1",
],
{ stdio: ["ignore", "pipe", "pipe"] },
);
// Track so cleanup() can terminate the merge during graceful shutdown.
activeProcesses.add(ffmpeg);
ffmpeg.once("exit", () => {
activeProcesses.delete(ffmpeg);
});
// Prevent the ffmpeg stderr from filling the pipe buffer / leaking.
let stderrBuf = "";
const MAX_STDERR = 4096;
ffmpeg.stderr?.on("data", (chunk: Buffer) => {
if (stderrBuf.length < MAX_STDERR) {
stderrBuf += chunk.toString("utf8");
}
});
ffmpeg.on("error", (err) => {
const msg =
err.message === "spawn ffmpeg ENOENT"
? "FFmpeg not found! Install ffmpeg in the container."
: err.message;
logger.error({ error: msg }, "Screen stream merge ffmpeg error");
});
ffmpeg.on("exit", (code) => {
const stderr = stderrBuf.trim();
logger.warn(
{ code, stderr: stderr.slice(-500) || undefined },
"Screen stream merge ffmpeg exited",
);
});
const stream = ffmpeg.stdout;
stream.setMaxListeners(32);
return stream;
}
/**
* Extract metadata (title, duration, thumbnail) from a media URL
* without downloading the audio stream.
@@ -1,3 +1,4 @@
import { PassThrough, type Readable } from "node:stream";
import type { Client } from "discord.js-selfbot-v13";
import { createChildLogger } from "@/shared/logger/index";
import {
@@ -54,6 +55,107 @@ export class ScreenShareController {
return this.active !== null;
}
/**
* Resolve the screen-share input with retry + first-byte validation.
*
* Transient YouTube 403s kill the merge ffmpeg BEFORE it produces any
* output; without validation the stream would "start" with a dead input
* and show a black tile forever. So after getDirectScreenInput resolves we
* tee the stream through a PassThrough and wait for the FIRST readable
* byte (or an error / early EOF). On failure the whole resolution is
* retried with a FRESH yt-dlp run (signed DASH URLs expire quickly the
* old URLs cannot simply be re-fetched).
*/
private async resolveInputWithRetry(source: string): Promise<Readable> {
const MAX_ATTEMPTS = 3;
let lastError: Error | null = null;
for (let attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) {
try {
const input = await getDirectScreenInput(source);
const tee = new PassThrough();
input.on("error", (err) => tee.destroy(err));
input.on("end", () => tee.end());
input.pipe(tee);
// If the merge process is stuck (no data, no exit) destroy the raw
// stream too so ffmpeg gets EPIPE on its next write and dies —
// otherwise every failed attempt leaks a merge process.
const destroyInput = () => {
try {
input.destroy();
} catch {
/* already gone */
}
};
await new Promise<void>((resolve, reject) => {
const timer = setTimeout(() => {
cleanup();
destroyInput();
tee.destroy(
new Error(
"Screen input produced no data within 12s — merge likely failed",
),
);
reject(
new Error(
"Screen input produced no data within 12s — merge likely failed",
),
);
}, 12000);
const onReadable = () => {
if (tee.readableLength > 0) {
cleanup();
resolve();
}
// readableLength === 0 can mean "EOF reached" — handled by onEnd.
};
const onError = (err: Error) => {
cleanup();
reject(err);
};
const onEnd = () => {
cleanup();
destroyInput();
reject(new Error("Screen input ended before producing any data"));
};
const cleanup = () => {
clearTimeout(timer);
tee.removeListener("readable", onReadable);
tee.removeListener("error", onError);
tee.removeListener("end", onEnd);
};
tee.once("readable", onReadable);
tee.once("error", onError);
tee.once("end", onEnd);
});
// Pass the tee onward — the encoder consumes the same buffered
// stream, so no data from the merge is lost.
return tee;
} catch (err) {
lastError = err instanceof Error ? err : new Error(String(err));
this.logger.warn(
{
attempt,
maxAttempts: MAX_ATTEMPTS,
error: lastError.message,
},
"Screen input resolution failed; retrying with fresh yt-dlp",
);
if (attempt < MAX_ATTEMPTS) {
await new Promise((r) => setTimeout(r, 1500 * attempt));
}
}
}
throw (
lastError ??
new Error("Screen input resolution failed after multiple attempts")
);
}
async start(source: string): Promise<ScreenSharePlayback> {
const status = this.getVoiceStatus();
if (!status.connected || !status.activeGuildId || !status.activeChannelId) {
@@ -65,7 +167,7 @@ export class ScreenShareController {
}
try {
const input = await getDirectScreenInput(source);
const input = await this.resolveInputWithRetry(source);
if (!this.streamer) {
this.streamer = new Streamer(this.client);
}
@@ -105,7 +207,11 @@ export class ScreenShareController {
frameRate: 30,
bitrateVideo: 2500,
bitrateVideoMax: 4000,
includeAudio: true,
// Video-only GoLive: the -f h264 output muxer cannot carry audio
// ("h264 muxer does not support any stream of type audio" → header
// write fails → empty stdout → demux 'Invalid data' → black tile).
// Audio is not delivered by the GoLive demux path anyway.
includeAudio: false,
videoCodec: normalizeVideoCodec("H264"),
});
const { command } = prepared;
@@ -145,6 +251,9 @@ export class ScreenShareController {
};
const done = playStream(prepared, this.streamer, {
type: "go-live",
width: 1280,
height: 720,
frameRate: 30,
})
.catch((err: unknown) => {
// Never let a stream failure become an unhandledRejection — that
@@ -56,6 +56,24 @@ export class VoiceTransmitter {
// Create PCM input stream
this.pcmStream = new PassThrough();
this.pcmStream.setMaxListeners(32); // drain listeners accumulate during backpressure
// Voice teardown (stop / disconnect / ffmpeg exit) destroys this stream
// while Redis PCM messages may still be in flight. Without a listener,
// EPIPE / ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END surface as
// an uncaughtException and crash the whole gateway.
this.pcmStream.on("error", (err: NodeJS.ErrnoException) => {
if (
err.code === "EPIPE" ||
err.code === "ERR_STREAM_DESTROYED" ||
err.code === "ERR_STREAM_WRITE_AFTER_END"
) {
logger.debug(
{ code: err.code },
"PCM stream closed during voice teardown — ignoring",
);
} else {
logger.error({ error: err.message }, "PCM stream error");
}
});
// Spawn FFmpeg to encode 24kHz mono PCM → OggOpus
// Input: 24kHz mono s16le (raw PCM)
@@ -146,7 +164,12 @@ export class VoiceTransmitter {
);
this.redisSub.on("message", (channel, message) => {
if (channel !== this.TRANSMIT_CHANNEL || !this.pcmStream) return;
if (
!this.isActive ||
channel !== this.TRANSMIT_CHANNEL ||
!this.pcmStream
)
return;
try {
const data = JSON.parse(message);
@@ -161,11 +184,21 @@ export class VoiceTransmitter {
this.draining = false;
// Re-acquire stream reference (could have been replaced by restart)
const currentStream = this.pcmStream;
if (!currentStream) return;
if (!currentStream || !this.isActive) return;
// Flush queued chunks
while (this.backpressureQueue.length > 0) {
const queued = this.backpressureQueue.shift()!;
if (!currentStream.write(queued)) break;
try {
if (!currentStream.write(queued)) break;
} catch (err) {
logger.debug(
{
error: err instanceof Error ? err.message : String(err),
},
"PCM flush write failed during teardown — ignoring",
);
break;
}
}
});
}
@@ -171,6 +171,24 @@ export const configSchema = z
.int()
.positive()
.default(30000),
// Term glossary — per-word Wikipedia lookups (via SearXNG) for words the
// LLM may not know (slang, jargon, regional language, foreign terms).
// Definitions are cached (in-memory + Redis) so repeat lookups are fast.
// Disable to skip glossary lookups entirely and analyze without them.
AI_GLOSSARY_ENABLED: z
.string()
.optional()
.transform((v) => v === "true")
.default(true),
// Max glossary terms looked up per analysis batch (keeps latency bounded).
AI_GLOSSARY_MAX_TERMS: z.coerce.number().int().min(1).max(20).default(6),
// Min word length for a term to be considered glossary-worthy.
AI_GLOSSARY_MIN_WORD_LENGTH: z.coerce
.number()
.int()
.min(2)
.max(20)
.default(5),
// ── AI Analysis Timing ──────────────────────────────────────────────
AI_ANALYSIS_DEBOUNCE_MS: z.coerce.number().positive().default(500),
@@ -22,11 +22,14 @@ describe("goLive port: codec + encoders", () => {
expect(normalizeVideoCodec("av1")).toBe("AV1");
});
it("software encoder exposes x264 libx264 superfast film", () => {
it("software encoder exposes x264 libx264 baseline zerolatency", () => {
const enc = Encoders.software()();
expect(enc.H264.name).toBe("libx264");
expect(enc.H264.options).toContain("-preset superfast");
expect(enc.H264.options).toContain("-tune film");
expect(enc.H264.options).toContain("-tune zerolatency");
// Baseline profile is REQUIRED to match the SDP's profile-level-id=42e01f
// (constrained baseline) — High-profile bitstreams fail to decode → black
expect(enc.H264.options).toContain("-profile:v baseline");
});
it("CodecPayloadType has opus + H264 entries", () => {
@@ -0,0 +1,130 @@
// Regression test: demux must emit frames from a LIVE stream that never
// ends (the NUT/H264 merge output during playback). The old implementation
// spooled the whole stream to a file first → deadlocked forever → 0 frames.
// v2: also validates ACCESS-UNIT grouping — each emitted frame must be a
// complete picture (parameter sets + slice), never a bare SPS/PPS/SEI NAL,
// and must be timestamped at the video frame rate (RTP +clockRate/fps).
// Run: npx tsx tests/golive-demux-live-e2e.ts [ffmpeg-path]
import { spawn } from "node:child_process";
import { PassThrough } from "node:stream";
import { demux } from "../src/goLive/Demuxer.js";
const FFMPEG = process.argv[2] ?? "ffmpeg";
// 1) Generate a 2s H264 test clip to a temp file
const clip = "/tmp/golive-live-test.h264";
await new Promise<void>((resolve, reject) => {
const p = spawn(
FFMPEG,
[
"-hide_banner", "-loglevel", "error",
"-f", "lavfi", "-i", "testsrc=size=640x360:rate=30:duration=2",
"-c:v", "libx264", "-preset", "ultrafast", "-pix_fmt", "yuv420p",
"-f", "h264", clip,
],
{ stdio: ["ignore", "ignore", "pipe"] },
);
let err = "";
p.stderr?.on("data", (d: Buffer) => (err += d.toString()));
p.on("close", (code) => (code === 0 ? resolve() : reject(new Error(err))));
});
// 2) Feed the clip through a PassThrough but DON'T end it (live semantics),
// with a small pause after the first chunk so demux has time to emit.
const input = new PassThrough();
const demuxPromise = demux(input, { format: "h264", frameRate: 30 });
const { video, close } = await demuxPromise;
if (!video) {
console.error("FAIL: demux returned no video stream");
process.exit(1);
}
interface Emitted {
data: Buffer;
duration: number;
timeBase: { num: number; den: number };
flags: number;
}
const frames: Emitted[] = [];
video.stream.on("data", (frame: Emitted) => {
frames.push(frame);
});
const fs = await import("node:fs");
const buf = fs.readFileSync(clip);
const chunkSize = 16384;
for (let i = 0; i < buf.length; i += chunkSize) {
input.write(buf.subarray(i, i + chunkSize));
if (i === 0) await new Promise((r) => setTimeout(r, 1500));
}
// Stream still open — if the old spool logic was here we'd never emit.
await new Promise((r) => setTimeout(r, 500));
// 3) Validate access-unit structure
const nalTypes = (frame: Buffer): number[] => {
const out: number[] = [];
let i = 0;
while (i < frame.length - 3) {
if (frame[i] === 0 && frame[i + 1] === 0 && frame[i + 2] === 1) {
const start = i;
let j = i + 3;
if (frame[j - 4] === 0 && j >= 4) {
// 4-byte start code already consumed by i pointing at the 3-byte tail
}
while (j < frame.length - 3) {
if (frame[j] === 0 && frame[j + 1] === 0 && frame[j + 2] === 1) break;
j++;
}
const nal = frame.subarray(start + 3, j);
if (nal.length > 0) out.push(nal[0] & 0x1f);
i = j;
} else {
i++;
}
}
return out;
};
let bareParamSetFrames = 0;
let framesWithoutSlice = 0;
let keyframesWithParamSets = 0;
let keyframesWithoutParamSets = 0;
for (const f of frames) {
const types = nalTypes(f.data);
const hasSlice = types.some((t) => t === 1 || t === 5);
const hasParams = types.some((t) => t === 7 || t === 8);
const isKey = (f.flags & 1) !== 0;
if (!hasSlice) framesWithoutSlice++;
if (types.length === 1 && (types[0] === 7 || types[0] === 8 || types[0] === 6)) {
bareParamSetFrames++;
}
if (isKey && hasParams) keyframesWithParamSets++;
if (isKey && !hasParams) keyframesWithoutParamSets++;
}
console.log(
`metadata: ${video.codecName} ${video.width}x${video.height} fps=${video.framerate_num}/${video.framerate_den}`,
);
console.log(`frames while stream OPEN (not ended): ${frames.length}`);
console.log(`frames w/o slice NAL: ${framesWithoutSlice}, bare param-set frames: ${bareParamSetFrames}`);
console.log(`keyframes with SPS/PPS: ${keyframesWithParamSets}, without: ${keyframesWithoutParamSets}`);
if (frames.length === 0) {
console.error("FAIL: no frames emitted while input still open (deadlock)");
close();
process.exit(1);
}
if (bareParamSetFrames > 0 || framesWithoutSlice > 0) {
console.error("FAIL: demux emitted bare parameter-set frames (must group into access units)");
close();
process.exit(1);
}
if (frames.some((f) => f.duration !== 1 || f.timeBase.den !== 30)) {
console.error("FAIL: frame duration/timeBase not 1/30 (RTP timestamp advance wrong)");
close();
process.exit(1);
}
input.end();
await new Promise((r) => setTimeout(r, 300));
close();
console.log("PASS: live stream demux works + access units grouped correctly");
process.exit(0);
@@ -1,16 +1,21 @@
// ═══════════════════════════════════════════════════════════════════════════════
// Screen share input resolution tests
//
// Verifies the decision logic of getDirectScreenInput:
// - merged progressive URL → returned directly
// - video+audio DASH pair → local ffmpeg merge (Readable)
// - neither → rejection
// getDirectScreenInput now streams the merged video+audio media straight from
// yt-dlp stdout (`-o -`) — same auth-handling mechanism as resolveMediaUrl for
// music. There is no manual URL fetch or local ffmpeg merge anymore.
//
// Both yt-dlp and ffmpeg are faked via PATH shim scripts so the test does not
// hit the network or need real binaries.
// yt-dlp is faked via a PATH shim script so the test does not hit the network
// or need real binaries.
// ═══════════════════════════════════════════════════════════════════════════════
import { chmodSync, mkdtempSync, rmSync, writeFileSync } from "node:fs";
import {
chmodSync,
mkdtempSync,
readFileSync,
rmSync,
writeFileSync,
} from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { Readable } from "node:stream";
@@ -25,28 +30,21 @@ const realPath = process.env.PATH;
beforeAll(() => {
fakeBinDir = mkdtempSync(join(tmpdir(), "gmw-fake-bins-"));
// Fake yt-dlp: prints the JSON file named in GMW_FAKE_YTDLP_JSON.
// If the file is missing → exits 1 (mimics yt-dlp failure).
// Fake yt-dlp: streams a few bytes to stdout (like `yt-dlp -o -` does).
// Modes (env):
// GMW_FAKE_YTDLP_FAIL=1 → stderr 403 + exit 8 WITHOUT stdout bytes
// (mimics a download rejected by YouTube).
const ytShim = `#!/usr/bin/env bash
if [ -n "$GMW_FAKE_YTDLP_JSON" ] && [ -f "$GMW_FAKE_YTDLP_JSON" ]; then
cat "$GMW_FAKE_YTDLP_JSON"
exit 0
if [ "$GMW_FAKE_YTDLP_FAIL" = "1" ]; then
echo "ERROR: [youtube] ...: 403 Forbidden (access denied)" >&2
exit 8
fi
echo "yt-dlp: fake JSON missing" >&2
exit 1
`;
writeFileSync(join(fakeBinDir, "yt-dlp"), ytShim);
chmodSync(join(fakeBinDir, "yt-dlp"), 0o755);
// Fake ffmpeg: writes a small nut-ish payload to stdout so the returned
// Readable actually emits data (the merge path in mergeScreenStreams).
const ffShim = `#!/usr/bin/env bash
# Fake ffmpeg ignore args, emit a few bytes so consumers see a live stream.
# Fake yt-dlp ignore args, emit a few bytes so consumers see a live stream.
head -c 4096 /dev/urandom
exit 0
`;
writeFileSync(join(fakeBinDir, "ffmpeg"), ffShim);
chmodSync(join(fakeBinDir, "ffmpeg"), 0o755);
writeFileSync(join(fakeBinDir, "yt-dlp"), ytShim);
chmodSync(join(fakeBinDir, "yt-dlp"), 0o755);
process.env.PATH = `${fakeBinDir}:${process.env.PATH}`;
});
@@ -59,84 +57,76 @@ afterAll(() => {
});
// ─── helpers ───────────────────────────────────────────────────────────────────
function writeFakeJson(payload: Record<string, unknown>): string {
const p = join(
tmpdir(),
`gmw-fake-${process.pid}-${Date.now()}-${Math.random().toString(36).slice(2)}.json`,
);
writeFileSync(p, JSON.stringify(payload));
return p;
}
function dashPairInfo(videoUrl: string, audioUrl: string) {
return {
url: null,
acodec: "none", // top-level is not a single merged format
vcodec: "av01",
requested_formats: [
{
format_id: "136",
vcodec: "avc1.4d401f",
acodec: "none",
url: videoUrl,
},
{ format_id: "140", vcodec: "none", acodec: "mp4a.40.2", url: audioUrl },
],
};
function consumeStream(stream: Readable): Promise<string> {
return new Promise<string>((resolve) => {
let got = 0;
stream.on("data", (chunk: Buffer) => {
got += chunk.length;
});
stream.on("error", () => resolve(`error-after-${got}B`));
stream.on("end", () => resolve(`end-after-${got}B`));
stream.resume();
});
}
// ─── tests ─────────────────────────────────────────────────────────────────────
describe("getDirectScreenInput", () => {
it("returns the single merged progressive URL when the info has one", async () => {
process.env.GMW_FAKE_YTDLP_JSON = writeFakeJson({
url: "https://cdn.example/progressive.mp4",
acodec: "mp4a.40.2",
vcodec: "avc1",
});
const result = await getDirectScreenInput("https://youtu.be/abc");
expect(result).toBe("https://cdn.example/progressive.mp4");
});
it("returns a live Readable when a video+audio DASH pair must be merged", async () => {
process.env.GMW_FAKE_YTDLP_JSON = writeFakeJson(
dashPairInfo(
"https://cdn.example/video.mp4",
"https://cdn.example/audio.m4a",
),
);
it("returns a live Readable and streams media bytes from yt-dlp stdout", async () => {
const result = await getDirectScreenInput("https://youtu.be/abc");
expect(Readable.isReadable(result)).toBe(true);
// The fake ffmpeg emits bytes; collect a chunk to prove the stream flows.
const bytes = await new Promise<number>((resolve, reject) => {
const stream = result as Readable;
let got = 0;
stream.on("data", (chunk: Buffer) => {
got += chunk.length;
});
stream.on("error", reject);
stream.on("end", () => resolve(got));
stream.resume();
});
expect(bytes).toBeGreaterThan(0);
const outcome = await consumeStream(result);
// The fake yt-dlp emits 4096 bytes → the stream must deliver them.
expect(outcome).toMatch(/^(error|end)-after-[1-9]\d*B$/);
});
it("rejects when yt-dlp returns neither a merged URL nor a format pair", async () => {
process.env.GMW_FAKE_YTDLP_JSON = writeFakeJson({
url: null,
acodec: "none",
vcodec: "none",
requested_formats: [],
});
await expect(getDirectScreenInput("https://youtu.be/abc")).rejects.toThrow(
/neither a merged progressive URL nor a video\+audio/,
);
it("destroys the stream with an error when yt-dlp fails before producing data (transient 403)", async () => {
// Simulate the production failure: yt-dlp's downloader hits a transient
// YouTube 403 and exits non-zero WITHOUT emitting a single byte. The
// returned Readable must terminate with zero bytes (error OR end) so the
// controller's resolveInputWithRetry retries with a fresh run instead of
// streaming a silent black tile.
process.env.GMW_FAKE_YTDLP_FAIL = "1";
try {
const result = await getDirectScreenInput("https://youtu.be/abc");
expect(Readable.isReadable(result)).toBe(true);
const outcome = await consumeStream(result);
expect(outcome).toMatch(/^(error|end)-after-0B$/);
} finally {
delete process.env.GMW_FAKE_YTDLP_FAIL;
}
});
it("rejects when yt-dlp exits non-zero", async () => {
process.env.GMW_FAKE_YTDLP_JSON = "/nonexistent/gmw-fake.json";
await expect(getDirectScreenInput("https://youtu.be/abc")).rejects.toThrow(
/screen input resolution exited with code 1/,
it("passes -o - (stdout streaming) and a temp dir to yt-dlp", async () => {
const argsDump = join(
tmpdir(),
`gmw-ytargs-${process.pid}-${Date.now()}.txt`,
);
process.env.GMW_FAKE_YTDLP_DUMP_ARGS = argsDump;
// Augment the fake to dump its argv.
const shim = `#!/usr/bin/env bash
printf '%s\\n' "$*" >> "$GMW_FAKE_YTDLP_DUMP_ARGS"
head -c 4096 /dev/urandom
exit 0
`;
const realPath2 = process.env.PATH;
const dir = fakeBinDir as unknown as string;
const existing = join(dir, "yt-dlp");
// Overwrite with the argv-dumping variant.
writeFileSync(existing, shim);
chmodSync(existing, 0o755);
try {
const result = await getDirectScreenInput("https://youtu.be/abc");
await consumeStream(result);
await new Promise((r) => setTimeout(r, 100));
const args = readFileSync(argsDump, "utf8").trim();
expect(args).toContain("-o -");
expect(args).toMatch(/gmw-ytdlp-/);
} finally {
delete process.env.GMW_FAKE_YTDLP_DUMP_ARGS;
rmSync(argsDump, { force: true });
process.env.PATH = realPath2;
}
});
});
@@ -0,0 +1,91 @@
// ═══════════════════════════════════════════════════════════════════════════
// Term glossary — pure extraction/formatting tests (no DB, Redis, or network)
// ═══════════════════════════════════════════════════════════════════════════
import { describe, expect, it } from "vitest";
import {
extractGlossaryTerms,
formatTermGlossary,
} from "../src/modules/ai-moderation/termGlossary.js";
describe("extractGlossaryTerms — filters out words the LLM already knows", () => {
it("returns [] for common conversational Indonesian", () => {
const terms = extractGlossaryTerms(
["anjay mabar yuk gaskeun gua gas", "iya bener banget sih"],
{ maxTerms: 6 },
);
expect(terms).toEqual([]);
});
it("extracts uncommon/foreign-looking words and skips stopwords + brands", () => {
const terms = extractGlossaryTerms(
[
"tadi gua baca soal tempeh di discord",
"kayaknya istilahnya shirkmaxxing deh",
],
{ maxTerms: 6 },
);
// "tempeh" and "shirkmaxxing" are candidates; "discord"/"istilahnya" are not
expect(terms).toContain("tempeh");
expect(terms).toContain("shirkmaxxing");
expect(terms).not.toContain("discord");
expect(terms).not.toContain("istilahnya");
});
it("strips URLs, mentions, and custom emoji before extracting", () => {
const terms = extractGlossaryTerms(
["cek https://example.com/foo <@123456> <:hadeh:987> kafircel"],
{ maxTerms: 6 },
);
expect(terms).toContain("kafircel");
expect(terms.some((t) => /example|hadeh|123/.test(t))).toBe(false);
});
it("extracts quoted phrases as a single term", () => {
const terms = extractGlossaryTerms(['dia bilang "kostum hewan" itu aneh'], {
maxTerms: 6,
});
expect(terms).toContain("kostum hewan");
});
it("skips repeated-char noise like wkwkwk and aaaaa", () => {
const terms = extractGlossaryTerms(["wkwkwkwk aaaaa xixixi"], {
maxTerms: 6,
});
expect(terms).toEqual([]);
});
it("respects maxTerms and prioritizes proper nouns", () => {
const terms = extractGlossaryTerms(
["aku suka Xenogears sama Chrono Cross terus Yakuza"],
{ maxTerms: 2 },
);
expect(terms.length).toBeLessThanOrEqual(2);
expect(terms[0]).toBe("Xenogears");
});
});
describe("formatTermGlossary — XML block shape", () => {
it("returns '' for an empty map", () => {
expect(formatTermGlossary(new Map())).toBe("");
});
it("wraps definitions in <term_glossary> with escaped attributes/content", () => {
const block = formatTermGlossary(
new Map([
[
"kafircel",
{
term: "kafircel",
definition: "sebutan <memes> untuk & orang",
sourceUrl: "https://id.wikipedia.org/wiki/Mem",
},
],
]),
);
expect(block).toContain("<term_glossary>");
expect(block).toContain('<term word="kafircel"');
expect(block).toContain("&lt;memes&gt;");
expect(block).toContain("&amp;");
expect(block).toContain("</term_glossary>");
});
});