feat(gateway): split video recording into silence-based segments like voice

Video (camera + screen share) DAVE stream-watch now produces per-burst
MP4 segments instead of one long .h264 per watch:
- Detects VIDEO_SILENCE_MS (4000ms) of no H264 packets → closes the
  current segment, muxes to MP4, registers in voice_recordings + uploads
  to TeleUploader, then reopens for the next burst (mirrors voice AfterSilence).
- Per-watch segment counter + per-segment depacketizer reset + closing
  guard + write-error swallow so races (silence close vs in-flight UDP
  packet) never corrupt files or crash the gateway.
- Frontend: recordings deck renders a native <video> player for MP4 rows
  (detected by filename), keeps single-playback registry across audio+video.
This commit is contained in:
asepharyana
2026-09-01 21:51:44 +07:00
committed by asepharyana
parent eb2f8f5b6f
commit 4ec9685194
5 changed files with 417 additions and 49 deletions
@@ -0,0 +1,33 @@
# Video Recording Splitting — Like Voice Recording
## Goal
Camera + screen share (stream watch) recording should split into per-burst
segments just like voice recording does — each time a streamer pauses/stops
and resumes, a new MP4 segment is created and registered in the DB + uploaded.
## Voice Recording Model (to replicate)
1. `receiver.speaking.start` → new OGG segment per burst
2. AfterSilence (4000ms) → stream "end" → segment finalized + uploaded
3. Each segment → DB insert → OGG→MP3 transcode → upload → update DB
4. File stored as `<userId>/<startTime>.ogg` + `.json`
## Video Recording Splitting
1. DAVE video RTP → depacketize H264 → write to current segment .h264
2. Silence detection: no H264 packets for 4000ms → close segment → flush →
mux to MP4 → insert DB record → upload → start new segment on next packet
3. Each segment: `<userId>/video-<channelId>-<startTime>.h264` → `.mp4`
4. DB: reuse `voice_recordings` table (filename indicates video, e.g. `video-XXX-1234.mp4`)
5. Upload: MP4 to TeleUploader (no transcode needed — MP4 plays everywhere)
## Files Modified
- `services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts`
— Main change: silence-based splitting + DB registration + upload
## Constants
- `VIDEO_SILENCE_MS = 4000` (matches voice AfterSilence)
- `VIDEO_MIN_SEGMENT_MS = 1000` (skip segments <1s — avoid noise)
## Verification
- `pnpm typecheck` in `services/discord-gateway`
- `pnpm build` (dist/ is the deployed artifact)
- Push → CI deploy → live test with a streamer
@@ -29,6 +29,7 @@ import {
promises as fsPromises,
type WriteStream,
} from "node:fs";
import { stat } from "node:fs/promises";
import path from "node:path";
import {
getVoiceConnection,
@@ -50,6 +51,19 @@ const HEADER_EXTENSION_BYTE = Buffer.from([190, 222]); // 0xBE 0xDE
/** Opus RTP payload type (120). Video arrives on non-opus payload types. */
const RTP_OPUS_PAYLOAD_TYPE = 120;
/**
* Video silence threshold — matches voice recording AfterSilence (4000ms).
* When no H264 video RTP packets arrive for this long, the current segment
* is closed and muxed to MP4 (like voice's "end of burst"). A new segment
* opens when packets resume.
*/
const VIDEO_SILENCE_MS = 4000;
/**
* Minimum segment duration — skip registering segments shorter than this
* (avoids creating tiny MP4 files from brief camera flashes).
*/
const VIDEO_MIN_SEGMENT_MS = 1000;
/** A watched user's live DAVE stream connection + file state. */
interface WatchState {
uid: string;
@@ -62,6 +76,12 @@ interface WatchState {
bytesWritten: number;
lastPacketAt: number;
startedAt: number;
/** Timestamp of the current segment's first packet (for silence detection). */
segmentStartAt: number;
/** Per-watch segment counter — makes each segment's filename unique. */
segmentSeq: number;
/** True while a segment close is in-flight (rejects new packet writes). */
closing: boolean;
_diag?: { video: number; maxLen: number; pts: Set<number> };
_diagNext?: number;
}
@@ -77,6 +97,37 @@ export function setStreamWatchRecordingsDir(dir: string): void {
_recordingsDir = dir;
}
/**
* Silence detection interval — runs every 2s and checks each active watch for
* VIDEO_SILENCE_MS of inactivity (no H264 packets). On silence, the current
* segment is closed, muxed to MP4, registered in the DB, and uploaded.
* Matches voice recording's AfterSilence behavior.
*/
const _silenceInterval = setInterval(() => {
const now = Date.now();
for (const [watchKey, watch] of watches) {
if (!watch.out || !watch.segmentStartAt) continue;
const silenceDuration = now - watch.lastPacketAt;
if (silenceDuration < VIDEO_SILENCE_MS) continue;
const segmentDuration = now - watch.segmentStartAt;
logger.info(
{
userId: watch.uid,
watchKey,
segmentDurationMs: segmentDuration,
bytes: watch.bytesWritten,
},
`Video silence detected (${silenceDuration}ms) — closing segment`,
);
// Fire-and-forget: close + finalize the segment.
void closeCurrentSegment(watch, watchKey);
}
}, 2_000);
// Prevent the interval from keeping the process alive.
if (_silenceInterval.unref) _silenceInterval.unref();
/** Test-only reset. */
export function __resetStreamWatchState(): void {
for (const [, w] of watches) closeWatch(w);
@@ -330,6 +381,9 @@ function connectWatch(
bytesWritten: 0,
lastPacketAt: Date.now(),
startedAt: Date.now(),
segmentStartAt: 0,
segmentSeq: 0,
closing: false,
});
}
@@ -530,6 +584,8 @@ function handleUdpMessage(
void openOutput(watchKey, watch, decrypted);
return;
}
// Reject writes while a segment close is in-flight (silence detected).
if (watch.closing) return;
watch.lastPacketAt = Date.now();
let nals: Buffer[];
try {
@@ -553,18 +609,203 @@ async function openOutput(
if (watch.out) return;
const dir = path.join(_recordingsDir, watch.uid);
await fsPromises.mkdir(dir, { recursive: true });
// Filename includes the per-watch segment counter so a new segment after
// silence never collides with the (possibly still in-flight) previous one.
const seq = watch.segmentSeq++;
const filePath = path.join(
dir,
`video-${watch.channelId}-${watch.startedAt}.h264`,
`video-${watch.channelId}-${watch.startedAt}-${seq}.h264`,
);
watch.filePath = filePath;
watch.out = createWriteStream(filePath);
// Swallow EPIPE / ERR_STREAM_WRITE_AFTER_END on the segment file when a
// close races an in-flight UDP write — an unhandled stream 'error' here
// would crash the gateway (same class as the media EPIPE incidents).
watch.out.on("error", () => {});
watch.segmentStartAt = Date.now();
logger.info(
{ userId: watch.uid, path: filePath },
"Video burst opened (DAVE video RTP received)",
"Video segment opened (DAVE video RTP received)",
);
}
/**
* Close the current segment's write stream, mux the H264 to MP4, and
* register the recording in the DB + upload. After this, the watch is
* ready for a new segment (silence ended → packets resume).
* Mirrors voice recording's finalizeSegment + uploadRecordingSegment flow.
*/
async function closeCurrentSegment(
watch: WatchState,
watchKey: string,
): Promise<void> {
if (!watch.out || watch.out.destroyed || watch.closing) return;
// Mark closing FIRST (synchronously) so the UDP handler stops writing.
watch.closing = true;
const filePath = watch.filePath;
const segmentStart = watch.segmentStartAt;
const bytes = watch.bytesWritten;
// Flush the write stream to disk.
const flushed = new Promise<void>((resolve) => {
if (!watch.out || watch.out.destroyed) {
resolve();
return;
}
watch.out.once("finish", () => resolve());
watch.out.end();
});
await flushed;
// Reset watch state for the next segment.
watch.out = undefined as unknown as WriteStream;
watch.filePath = "";
watch.bytesWritten = 0;
watch.segmentStartAt = 0;
watch.closing = false;
// Reset the depacketizer so a partial FU-A fragment from this segment
// doesn't bleed into the next segment (it only writes complete NALs).
watch.depacketizer.reset();
if (!filePath) return;
// Discard very short segments (<1s) — avoids creating tiny MP4 files.
const durationMs = Date.now() - segmentStart;
if (durationMs < VIDEO_MIN_SEGMENT_MS) {
logger.debug(
{ userId: watch.uid, durationMs },
"Video segment too short, discarding",
);
fsPromises.unlink(filePath).catch(() => {});
return;
}
// Mux H264 to MP4 and register in DB + upload.
try {
const mp4 = await muxToMp4(filePath);
await finalizeVideoSegment({
mp4Path: mp4,
userId: watch.uid,
guildId: watch.guildId,
channelId: watch.channelId,
startTime: segmentStart,
durationMs,
bytes,
});
logger.info(
{ userId: watch.uid, mp4, durationMs, bytes },
"Video segment finalized + uploaded",
);
} catch (err) {
logger.warn(
{
userId: watch.uid,
err: err instanceof Error ? err.message : String(err),
},
"Video segment finalize failed (raw .h264 kept)",
);
}
}
/**
* Register a completed video MP4 segment in the DB and upload it.
* Reuses the voice_recordings table (filename indicates video).
*/
async function finalizeVideoSegment(input: {
mp4Path: string;
userId: string;
guildId: string;
channelId: string;
startTime: number;
durationMs: number;
bytes: number;
}): Promise<void> {
const { mp4Path, userId, guildId, channelId, startTime, durationMs, bytes } =
input;
const segmentId = `video-${userId}-${startTime}`;
const fileName = path.basename(mp4Path);
try {
// Get file size and register in DB.
const fileStats = await stat(mp4Path);
const { insertVoiceRecording } = await import(
"../../shared/database/voiceRecordingRepo.js"
);
await insertVoiceRecording({
id: segmentId,
user_id: userId,
username: userId, // Will be enriched by frontend from user_profiles
avatar_url: null,
guild_id: guildId,
channel_id: channelId,
channel_name: null,
filename: fileName,
size_bytes: fileStats.size,
upload_status: "pending",
created_at: Date.now(),
});
// Upload MP4 to TeleUploader.
const { uploadToTele } = await import("../../shared/uploader.js");
const { config } = await import("../../shared/config/config.js");
const fileBuffer = await fsPromises.readFile(mp4Path);
const uploadResult = await uploadToTele({
buffer: fileBuffer,
filename: fileName,
contentType: "video/mp4",
uploadUrl: config.TELE_UPLOAD_URL,
retries: 3,
});
// Update DB with upload URL.
const { updateVoiceRecordingAsUploaded } = await import(
"../../shared/database/voiceRecordingRepo.js"
);
await updateVoiceRecordingAsUploaded(
segmentId,
uploadResult.url,
Date.now(),
);
logger.info(
{ segmentId, url: uploadResult.url, durationMs, bytes },
"Video segment uploaded successfully",
);
// Broadcast via EventBroadcaster.
const { _eventBroadcaster } = await import("./recorder.js");
if (_eventBroadcaster) {
_eventBroadcaster
.voiceRecordingUploaded({
id: segmentId,
user_id: userId,
username: userId,
avatar_url: null,
guild_id: guildId,
channel_id: channelId,
channel_name: null,
filename: fileName,
size_bytes: fileStats.size,
download_url: uploadResult.url,
upload_status: "uploaded",
created_at: Date.now(),
uploaded_at: Date.now(),
})
.catch(() => {});
}
} catch (err) {
const errorMsg = err instanceof Error ? err.message : String(err);
logger.error(
{ segmentId, error: errorMsg },
"Failed to upload video segment",
);
// Mark as failed in DB (best-effort).
const { updateVoiceRecordingAsFailed } = await import(
"../../shared/database/voiceRecordingRepo.js"
);
await updateVoiceRecordingAsFailed(segmentId, errorMsg).catch(() => {});
}
}
function closeWatch(watch: WatchState): void {
if (!watch) return;
try {
@@ -572,28 +813,9 @@ function closeWatch(watch: WatchState): void {
} catch {
/* ignore */
}
// Finalize any open segment (mux to MP4 + DB + upload).
if (watch.out && !watch.out.destroyed) {
const flushed = new Promise<void>((resolve) => {
watch.out.once("finish", () => resolve());
watch.out.end();
});
void flushed
.then(() => muxToMp4(watch.filePath))
.then((mp4) =>
logger.info(
{ userId: watch.uid, mp4, bytes: watch.bytesWritten },
"Video muxed to mp4",
),
)
.catch((err) =>
logger.warn(
{
userId: watch.uid,
err: err instanceof Error ? err.message : String(err),
},
"Video mux failed (raw .h264 kept)",
),
);
void closeCurrentSegment(watch, `${watch.guildId}:${watch.uid}`);
}
}
@@ -55,6 +55,7 @@ import type {
SpeakerSummary,
VoiceRecording,
} from "@/lib/types";
import { isVideoRecording } from "@/lib/types/recording";
import { useWebSocket } from "@/lib/ws/context";
const ALL = "__all__";
@@ -557,7 +558,10 @@ export function RecordingsView({
{r.download_url ? (
<RecordingAudioPlayer
src={r.download_url}
label={`Voice recording by ${r.username}`}
label={`${
isVideoRecording(r) ? "Video" : "Voice"
} recording by ${r.username}`}
video={isVideoRecording(r)}
onPlayStateChange={(active) =>
setPlayingId((prev) => {
if (active) return r.id;
@@ -21,9 +21,45 @@ function formatTime(sec: number): string {
return `${m}:${s.toString().padStart(2, "0")}`;
}
/** Shared event wiring for a video element (kept small to avoid dup logic). */
function wireVideoEvents(
el: HTMLVideoElement,
onLoadedMeta: () => void,
onTime: () => void,
onEnd: () => void,
onPause: () => void,
onPlaying: () => void,
): void {
el.addEventListener("loadedmetadata", onLoadedMeta);
el.addEventListener("durationchange", onLoadedMeta);
el.addEventListener("timeupdate", onTime);
el.addEventListener("ended", onEnd);
el.addEventListener("pause", onPause);
el.addEventListener("playing", onPlaying);
}
/** Remove the listeners added by {@link wireVideoEvents}. */
function unwireVideoEvents(
el: HTMLVideoElement,
onLoadedMeta: () => void,
onTime: () => void,
onEnd: () => void,
onPause: () => void,
onPlaying: () => void,
): void {
el.removeEventListener("loadedmetadata", onLoadedMeta);
el.removeEventListener("durationchange", onLoadedMeta);
el.removeEventListener("timeupdate", onTime);
el.removeEventListener("ended", onEnd);
el.removeEventListener("pause", onPause);
el.removeEventListener("playing", onPlaying);
}
interface Props {
src: string;
label?: string;
/** True when the recording is video (camera/screen-share MP4). */
video?: boolean;
/** Lifted state: parent highlights the card that owns the active player. */
onPlayStateChange?: (playing: boolean) => void;
className?: string;
@@ -34,20 +70,28 @@ interface Props {
* play/pause with buffering spinner, click-to-seek progress bar, time label,
* animated equalizer bars while playing, and single-playback enforcement
* (starting one clip pauses all others).
*
* For `video` recordings it renders a native `<video controls>` instead —
* video needs the browser's own scrubbing/fullscreen UI — but still honors
* the single-playback registry and lifts play state to the parent card.
*/
export function RecordingAudioPlayer({
src,
label = "Voice recording",
video = false,
onPlayStateChange,
className,
}: Props) {
const videoRef = useRef<HTMLVideoElement | null>(null);
const audioRef = useRef<HTMLAudioElement | null>(null);
const [playing, setPlaying] = useState(false);
const [buffering, setBuffering] = useState(false);
const [current, setCurrent] = useState(0);
const [duration, setDuration] = useState(0);
// Audio: hide an HTMLAudioElement and drive it with the custom scrubber.
useEffect(() => {
if (video) return;
const audio = new Audio();
audio.preload = "metadata";
audio.src = src;
@@ -107,37 +151,100 @@ export function RecordingAudioPlayer({
audio.src = "";
audioRef.current = null;
};
}, [src]);
}, [src, video]);
// Video: wire the rendered <video controls> for single-playback + lifted state.
useEffect(() => {
if (!video) return;
const el = videoRef.current;
if (!el) return;
const onLoadedMeta = () => setDuration(el.duration || 0);
const onTime = () => setCurrent(el.currentTime);
const onEnd = () => {
setPlaying(false);
setCurrent(0);
};
const onPause = () => setPlaying(false);
const onPlaying = () => setPlaying(true);
wireVideoEvents(el, onLoadedMeta, onTime, onEnd, onPause, onPlaying);
// Single playback: pausing any other video/audio that starts.
const pauseThis = () => el.pause();
let unregister: (() => void) | null = null;
const onPlayEvt = () => {
for (const other of activePlayers) {
if (other !== pauseThis) other();
}
unregister?.();
unregister = registerPlayer(pauseThis);
};
el.addEventListener("play", onPlayEvt);
return () => {
unregister?.();
unwireVideoEvents(el, onLoadedMeta, onTime, onEnd, onPause, onPlaying);
el.removeEventListener("play", onPlayEvt);
};
}, [video]);
useEffect(() => {
onPlayStateChange?.(playing || buffering);
}, [playing, buffering, onPlayStateChange]);
const toggle = useCallback(() => {
const audio = audioRef.current;
if (!audio) return;
if (audio.paused) {
const el = video ? videoRef.current : audioRef.current;
if (!el) return;
if (el.paused) {
setBuffering(true);
void audio.play().catch(() => setBuffering(false));
void el.play().catch(() => setBuffering(false));
} else {
audio.pause();
el.pause();
}
}, []);
}, [video]);
const seek = useCallback((e: React.MouseEvent<HTMLDivElement>) => {
const audio = audioRef.current;
if (!audio || !Number.isFinite(audio.duration)) return;
const rect = e.currentTarget.getBoundingClientRect();
const ratio = Math.min(
1,
Math.max(0, (e.clientX - rect.left) / rect.width),
);
audio.currentTime = ratio * audio.duration;
setCurrent(audio.currentTime);
}, []);
const seek = useCallback(
(e: React.MouseEvent<HTMLDivElement>) => {
const el = video ? videoRef.current : audioRef.current;
if (!el || !Number.isFinite(el.duration)) return;
const rect = e.currentTarget.getBoundingClientRect();
const ratio = Math.min(
1,
Math.max(0, (e.clientX - rect.left) / rect.width),
);
el.currentTime = ratio * el.duration;
setCurrent(el.currentTime);
},
[video],
);
const pct = duration > 0 ? (current / duration) * 100 : 0;
if (video) {
return (
<div
className={cn(
"rounded-[8px] border bg-surface-2 p-1.5 transition-colors",
playing
? "border-signal/40 shadow-[0_0_24px_-10px_var(--color-signal-glow)]"
: "border-hairline",
className,
)}
role="group"
aria-label={label}
>
<video
ref={videoRef}
src={src}
controls
playsInline
preload="metadata"
className="aspect-video w-full rounded-[5px] bg-black object-contain"
/>
</div>
);
}
return (
<div
className={cn(
@@ -183,15 +290,12 @@ export function RecordingAudioPlayer({
tabIndex={0}
onClick={seek}
onKeyDown={(e) => {
const audio = audioRef.current;
if (!audio || !Number.isFinite(audio.duration)) return;
const el = video ? videoRef.current : audioRef.current;
if (!el || !Number.isFinite(el.duration)) return;
if (e.key === "ArrowRight")
audio.currentTime = Math.min(
audio.duration,
audio.currentTime + 5,
);
el.currentTime = Math.min(el.duration, el.currentTime + 5);
if (e.key === "ArrowLeft")
audio.currentTime = Math.max(0, audio.currentTime - 5);
el.currentTime = Math.max(0, el.currentTime - 5);
}}
className="group relative h-4 cursor-pointer"
>
@@ -16,6 +16,11 @@ export interface VoiceRecording {
transcription?: string | null;
}
/** True when the recording file is video (camera/screen-share MP4). */
export function isVideoRecording(r: Pick<VoiceRecording, "filename">): boolean {
return /\.(mp4|webm|mov|h264)$/i.test(r.filename);
}
export interface PaginatedRecordings {
items: VoiceRecording[];
nextCursor: string | null;