diff --git a/services/discord-gateway/src/goLive/Demuxer.ts b/services/discord-gateway/src/goLive/Demuxer.ts index a2b5c9d..1ae1fe2 100644 --- a/services/discord-gateway/src/goLive/Demuxer.ts +++ b/services/discord-gateway/src/goLive/Demuxer.ts @@ -17,6 +17,9 @@ 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 @@ -145,7 +148,7 @@ export async function probeStreams( */ export async function demux( input: string | PassThrough, - _opts: { format: string }, + opts: { format: string; frameRate?: number }, ): Promise<{ video: DemuxedStream | undefined; audio: DemuxedStream | undefined; @@ -283,25 +286,55 @@ export async function demux( new Promise((resolve) => setTimeout(resolve, 1500)), ]); - // 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. + // 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=${nal.length}B key=${isKeyFrame}`, + `[goLive:Demuxer] frames=${frameCount} last=${au.length}B key=${isKey}`, ); } }; @@ -350,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 @@ -369,11 +415,7 @@ 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(); }); diff --git a/services/discord-gateway/src/goLive/prepareStream.ts b/services/discord-gateway/src/goLive/prepareStream.ts index e32e790..f29d92f 100644 --- a/services/discord-gateway/src/goLive/prepareStream.ts +++ b/services/discord-gateway/src/goLive/prepareStream.ts @@ -266,6 +266,8 @@ export async function playStream( 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}`, diff --git a/services/discord-gateway/tests/golive-demux-live-e2e.ts b/services/discord-gateway/tests/golive-demux-live-e2e.ts index 613f9de..20da726 100644 --- a/services/discord-gateway/tests/golive-demux-live-e2e.ts +++ b/services/discord-gateway/tests/golive-demux-live-e2e.ts @@ -1,6 +1,9 @@ // 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"; @@ -29,14 +32,22 @@ await new Promise((resolve, reject) => { // 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" }); +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); +} -let frames = 0; -let bytes = 0; -video.stream.on("data", (frame: { data: Buffer }) => { - frames++; - bytes += frame.data.length; +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"); @@ -49,15 +60,71 @@ for (let i = 0; i < buf.length; i += chunkSize) { // Stream still open — if the old spool logic was here we'd never emit. await new Promise((r) => setTimeout(r, 500)); -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}, bytes: ${bytes}`); -if (frames === 0) { +// 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"); +console.log("PASS: live stream demux works + access units grouped correctly"); process.exit(0);