diff --git a/.hermes/plans/2026-09-01_video-recording-splitting.md b/.hermes/plans/2026-09-01_video-recording-splitting.md new file mode 100644 index 00000000..cfced0d1 --- /dev/null +++ b/.hermes/plans/2026-09-01_video-recording-splitting.md @@ -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 `/.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: `/video--.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 diff --git a/services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts b/services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts index 0ea1cec1..f50da660 100644 --- a/services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts +++ b/services/discord-gateway/src/modules/voice-recording/streamWatchReceiver.ts @@ -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 }; _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 { + 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((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 { + 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((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}`); } } diff --git a/services/frontend/src/app/(dashboard)/recordings/view.tsx b/services/frontend/src/app/(dashboard)/recordings/view.tsx index b2013603..7cf7c7cd 100644 --- a/services/frontend/src/app/(dashboard)/recordings/view.tsx +++ b/services/frontend/src/app/(dashboard)/recordings/view.tsx @@ -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 ? ( setPlayingId((prev) => { if (active) return r.id; diff --git a/services/frontend/src/components/voice/recording-audio-player.tsx b/services/frontend/src/components/voice/recording-audio-player.tsx index ada9cce8..9afbc09c 100644 --- a/services/frontend/src/components/voice/recording-audio-player.tsx +++ b/services/frontend/src/components/voice/recording-audio-player.tsx @@ -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 `