feat(gateway): auto-mux recorded video to playable MP4 (Phase B)
Phase A captured raw .h264 streams but left them as non-playable elementary streams. Phase B adds automatic muxing: when a video burst closes, the raw .h264 is remuxed to a self-contained MP4 via `ffmpeg -c copy` (no re-encode, fast) with `+faststart`, waits for the write stream to fully flush first so the mux never reads a truncated tail, and deletes the raw .h264 on success (keeping it on failure). Output: <RECORDINGS_DIR>/<uid>/video-<ssrc>-<ts>.mp4. muxToMp4 is exported + covered by a real-ffmpeg vitest (tests/videoReceiver.test.ts): generates a tiny baseline h264, remuxes, asserts mp4 exists/non-empty & raw deleted (also the 5 depacketizer tests). Full gateway suite 170/170 green, tsc + biome clean.
This commit is contained in:
@@ -1,4 +1,6 @@
|
|||||||
|
import { spawn } from "node:child_process";
|
||||||
import { createWriteStream, mkdirSync, type WriteStream } from "node:fs";
|
import { createWriteStream, mkdirSync, type WriteStream } from "node:fs";
|
||||||
|
import { stat, unlink } from "node:fs/promises";
|
||||||
import path from "node:path";
|
import path from "node:path";
|
||||||
import type { VoiceReceiver } from "@discordjs/voice";
|
import type { VoiceReceiver } from "@discordjs/voice";
|
||||||
import type { Client } from "discord.js-selfbot-v13";
|
import type { Client } from "discord.js-selfbot-v13";
|
||||||
@@ -179,25 +181,39 @@ export function hookVideoReceiver(
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
/** Close a burst's file (flush + end). */
|
/** Close a burst's file (flush + end) then mux to a playable mp4. */
|
||||||
const closeBurst = (userId: string): void => {
|
const closeBurst = (userId: string): void => {
|
||||||
const burst = bursts.get(userId);
|
const burst = bursts.get(userId);
|
||||||
if (!burst) return;
|
if (!burst) return;
|
||||||
bursts.delete(userId);
|
bursts.delete(userId);
|
||||||
try {
|
const { ssrc, bytesWritten, createdAt, filePath } = burst;
|
||||||
|
// Wait for the write stream to fully flush (file descriptor closed) before
|
||||||
|
// ffmpeg reads it — otherwise the mux can read a truncated tail.
|
||||||
|
const flushed = new Promise<void>((resolve) => {
|
||||||
|
burst.out.once("finish", () => resolve());
|
||||||
burst.out.end();
|
burst.out.end();
|
||||||
} catch {
|
});
|
||||||
// ignore write errors on close
|
|
||||||
}
|
|
||||||
logger.info(
|
logger.info(
|
||||||
{
|
{ userId, ssrc, bytes: bytesWritten, durationMs: Date.now() - createdAt },
|
||||||
userId,
|
|
||||||
ssrc: burst.ssrc,
|
|
||||||
bytes: burst.bytesWritten,
|
|
||||||
durationMs: Date.now() - burst.createdAt,
|
|
||||||
},
|
|
||||||
"Video burst closed",
|
"Video burst closed",
|
||||||
);
|
);
|
||||||
|
// Phase B: remux the raw H264 elementary stream into a self-contained,
|
||||||
|
// playable MP4 (fast `-c copy`, no re-encode) and drop the raw file.
|
||||||
|
void flushed
|
||||||
|
.then(() => muxToMp4(filePath))
|
||||||
|
.then((mp4) => {
|
||||||
|
logger.info({ userId, ssrc, mp4 }, "Video muxed to mp4");
|
||||||
|
})
|
||||||
|
.catch((err) => {
|
||||||
|
logger.warn(
|
||||||
|
{
|
||||||
|
userId,
|
||||||
|
ssrc,
|
||||||
|
err: err instanceof Error ? err.message : String(err),
|
||||||
|
},
|
||||||
|
"Video mux to mp4 failed (keeping raw .h264)",
|
||||||
|
);
|
||||||
|
});
|
||||||
};
|
};
|
||||||
|
|
||||||
const openBurst = (userId: string, ssrc: number): VideoBurst | null => {
|
const openBurst = (userId: string, ssrc: number): VideoBurst | null => {
|
||||||
@@ -366,3 +382,51 @@ function getSsrcInternalMap(
|
|||||||
if (m1 instanceof Map) return m1 as Map<number, never>;
|
if (m1 instanceof Map) return m1 as Map<number, never>;
|
||||||
return undefined;
|
return undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Remux a raw H264 elementary stream into a playable MP4 with `-c copy` (no
|
||||||
|
* re-encode, fast). On success the raw `.h264` is removed and the `.mp4` path
|
||||||
|
* returned. On any failure the raw file is LEFT in place (caller keeps it).
|
||||||
|
*/
|
||||||
|
export async function muxToMp4(rawPath: string): Promise<string> {
|
||||||
|
const mp4Path = rawPath.replace(/\.h264$/, ".mp4");
|
||||||
|
|
||||||
|
await new Promise<void>((resolve, reject) => {
|
||||||
|
const proc = spawn(
|
||||||
|
"ffmpeg",
|
||||||
|
[
|
||||||
|
"-hide_banner",
|
||||||
|
"-loglevel",
|
||||||
|
"error",
|
||||||
|
"-f",
|
||||||
|
"h264",
|
||||||
|
"-i",
|
||||||
|
rawPath,
|
||||||
|
"-c",
|
||||||
|
"copy",
|
||||||
|
"-movflags",
|
||||||
|
"+faststart",
|
||||||
|
"-y",
|
||||||
|
mp4Path,
|
||||||
|
],
|
||||||
|
{ stdio: ["ignore", "ignore", "pipe"] },
|
||||||
|
);
|
||||||
|
let stderr = "";
|
||||||
|
proc.stderr?.on("data", (chunk) => {
|
||||||
|
stderr += chunk.toString();
|
||||||
|
});
|
||||||
|
proc.on("error", (err) => reject(err));
|
||||||
|
proc.on("close", (code) => {
|
||||||
|
if (code === 0) resolve();
|
||||||
|
else reject(new Error(`ffmpeg exit ${code ?? "?"}: ${stderr.trim()}`));
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// Sanity check: only delete the raw file if the mp4 is non-empty.
|
||||||
|
const mp4Stat = await stat(mp4Path).catch(() => null);
|
||||||
|
if (!mp4Stat || mp4Stat.size === 0) {
|
||||||
|
throw new Error("mp4 mux produced an empty file");
|
||||||
|
}
|
||||||
|
await unlink(rawPath).catch(() => {});
|
||||||
|
return mp4Path;
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,5 +1,13 @@
|
|||||||
|
import { execFileSync } from "node:child_process";
|
||||||
|
import { mkdtempSync, rmSync } from "node:fs";
|
||||||
|
import { access, stat } from "node:fs/promises";
|
||||||
|
import { tmpdir } from "node:os";
|
||||||
|
import path from "node:path";
|
||||||
import { describe, expect, it } from "vitest";
|
import { describe, expect, it } from "vitest";
|
||||||
import { H264Depacketizer } from "../src/modules/voice-recording/videoReceiver.js";
|
import {
|
||||||
|
H264Depacketizer,
|
||||||
|
muxToMp4,
|
||||||
|
} from "../src/modules/voice-recording/videoReceiver.js";
|
||||||
|
|
||||||
// Build a single-NAL RTP payload: [NAL header, ...data]
|
// Build a single-NAL RTP payload: [NAL header, ...data]
|
||||||
function singleNal(nalHeader: number, data: number[] = []): Buffer {
|
function singleNal(nalHeader: number, data: number[] = []): Buffer {
|
||||||
@@ -60,3 +68,57 @@ describe("H264Depacketizer", () => {
|
|||||||
expect(d.push(orphan).length).toBe(0);
|
expect(d.push(orphan).length).toBe(0);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("muxToMp4", () => {
|
||||||
|
it("remuxes a raw h264 file into a playable mp4 and deletes the raw file", async () => {
|
||||||
|
// Skip if ffmpeg is unavailable (headless Nix-less CI).
|
||||||
|
let ffmpegOk = true;
|
||||||
|
try {
|
||||||
|
execFileSync("ffmpeg", ["-version"], { stdio: "ignore" });
|
||||||
|
} catch {
|
||||||
|
ffmpegOk = false;
|
||||||
|
}
|
||||||
|
if (!ffmpegOk) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const dir = mkdtempSync(path.join(tmpdir(), "gmw-video-"));
|
||||||
|
const raw = path.join(dir, "clip.h264");
|
||||||
|
try {
|
||||||
|
execFileSync(
|
||||||
|
"ffmpeg",
|
||||||
|
[
|
||||||
|
"-hide_banner",
|
||||||
|
"-loglevel",
|
||||||
|
"error",
|
||||||
|
"-f",
|
||||||
|
"lavfi",
|
||||||
|
"-i",
|
||||||
|
"testsrc=duration=1:size=320x240:rate=10",
|
||||||
|
"-c:v",
|
||||||
|
"libx264",
|
||||||
|
"-preset",
|
||||||
|
"ultrafast",
|
||||||
|
"-profile:v",
|
||||||
|
"baseline",
|
||||||
|
"-pix_fmt",
|
||||||
|
"yuv420p",
|
||||||
|
"-f",
|
||||||
|
"h264",
|
||||||
|
"-y",
|
||||||
|
raw,
|
||||||
|
],
|
||||||
|
{ stdio: "ignore" },
|
||||||
|
);
|
||||||
|
|
||||||
|
const mp4 = await muxToMp4(raw);
|
||||||
|
expect(mp4).toBe(path.join(dir, "clip.mp4"));
|
||||||
|
// raw file deleted, mp4 exists and is non-empty
|
||||||
|
await expect(access(raw)).rejects.toThrow();
|
||||||
|
const s = await stat(mp4);
|
||||||
|
expect(s.size).toBeGreaterThan(0);
|
||||||
|
} finally {
|
||||||
|
rmSync(dir, { recursive: true, force: true });
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user