// NextNVR v0.3.3 — FFmpeg stream recorder // Records 5-minute segments with rec_ prefix and second-level timestamps. // Output: /mnt/recordings/{cam-name}/rec_{YYYY-MM-DD-HH-MM-SS}.mp4 package main import ( "fmt" "log" "os" "os/exec" "path/filepath" "strings" "sync" "time" ) // RecorderManager manages FFmpeg recording processes for all cameras. type RecorderManager struct { mu sync.Mutex processes map[string]*recorderProcess config StorageConfig } type recorderProcess struct { camera CameraConfig cmd *exec.Cmd done chan struct{} } // NewRecorderManager creates a new recorder manager. func NewRecorderManager(cfg StorageConfig) *RecorderManager { return &RecorderManager{ processes: make(map[string]*recorderProcess), config: cfg, } } // StartAll launches recording goroutines for all enabled cameras. // Cameras are started with a 500ms stagger to avoid I/O and CPU spikes. func (rm *RecorderManager) StartAll(cameras []CameraConfig) { for _, cam := range cameras { if !cam.Enabled || !cam.Record { log.Printf("recorder: skipping camera %s (enabled=%v, record=%v)", cam.ID, cam.Enabled, cam.Record) continue } go rm.startRecorder(cam) time.Sleep(500 * time.Millisecond) // staggered startup } log.Printf("recorder: all cameras launched") } // StopAll terminates all recording processes. func (rm *RecorderManager) StopAll() { rm.mu.Lock() defer rm.mu.Unlock() for id, proc := range rm.processes { log.Printf("recorder: stopping camera %s", id) close(proc.done) if proc.cmd != nil && proc.cmd.Process != nil { proc.cmd.Process.Signal(os.Interrupt) } } log.Println("recorder: all processes stopped") } // startRecorder runs the per-segment recording loop for a single camera. // // Architecture: // 1. Build RTSP URL from camera config. // 2. Create output directory: /mnt/recordings/{cam-name}/ // 3. LOOP: // a. Compute filename: YYYY-MM-DD-HH-MM.part.mp4 (start-of-segment timestamp) // b. Launch ffmpeg -t 300 -i {rtsp} -c copy → .part.mp4 // c. On success: atomic rename .part.mp4 → .mp4 // d. On failure: .part file stays (picked up by cleaner), wait 5s, retry // // Each segment is exactly 5 minutes. Gaps between segments are ~1-2s // (FFmpeg restart overhead). The -t 300 flag provides a hard timeout // preventing hung FFmpeg processes. func (rm *RecorderManager) startRecorder(cam CameraConfig) { rtspURL := cam.RTSPMain if rtspURL == "" { rtspURL = fmt.Sprintf("rtsp://%s:%s@%s:554/Streaming/Channels/101", cam.Username, cam.Password, cam.IP) } log.Printf("recorder: starting %s → %s", cam.ID, rtspURL) rm.mu.Lock() proc := &recorderProcess{ camera: cam, done: make(chan struct{}), } rm.processes[cam.ID] = proc rm.mu.Unlock() // Use camera name for the directory (human-readable). // Fall back to cam ID if name is empty. camDirName := cam.Name if camDirName == "" { camDirName = cam.ID } outDir := filepath.Join(rm.config.RecordingsPath, camDirName) if err := os.MkdirAll(outDir, 0755); err != nil { log.Printf("recorder: %s — mkdir %s failed: %v", cam.ID, outDir, err) return } for { select { case <-proc.done: log.Printf("recorder: %s — shutdown signal received", cam.ID) return default: } // Generate filename: rec_{YYYY-MM-DD-HH-MM-SS}.mp4 now := time.Now() timeStr := now.Format("2006-01-02-15-04-05") partFile := filepath.Join(outDir, "rec_"+timeStr+".part.mp4") finalFile := filepath.Join(outDir, "rec_"+timeStr+".mp4") // Build FFmpeg command with hard 5-minute timeout. cmd := exec.Command("ffmpeg", "-hide_banner", "-loglevel", "error", "-rtsp_transport", "tcp", "-i", rtspURL, "-an", // drop audio (pcm_mulaw not supported in MP4) "-c:v", "copy", // stream-copy video: zero transcoding, near-zero CPU "-t", "300", // hard stop after 300 seconds (5 minutes) "-y", // overwrite output without asking partFile, ) // Silence FFmpeg by default; uncomment for debugging. // cmd.Stdout = os.Stdout // cmd.Stderr = os.Stderr startTime := time.Now() if err := cmd.Start(); err != nil { log.Printf("recorder: %s — FFmpeg start failed: %v", cam.ID, err) time.Sleep(5 * time.Second) continue } proc.cmd = cmd log.Printf("recorder: %s — segment %s (PID %d)", cam.ID, timeStr, cmd.Process.Pid) err := cmd.Wait() proc.cmd = nil elapsed := time.Since(startTime) if err != nil { // FFmpeg exited with error or was killed. // .part file remains — will be cleaned up later. if strings.Contains(err.Error(), "signal: killed") || strings.Contains(err.Error(), "signal: interrupt") { log.Printf("recorder: %s — segment %s killed after %v (shutdown)", cam.ID, timeStr, elapsed.Round(time.Second)) return } log.Printf("recorder: %s — segment %s failed after %v: %v — retrying in 5s", cam.ID, timeStr, elapsed.Round(time.Second), err) // Remove partial file on error so it doesn't accumulate. os.Remove(partFile) time.Sleep(5 * time.Second) continue } // Segment completed successfully — atomic rename. if err := os.Rename(partFile, finalFile); err != nil { log.Printf("recorder: %s — rename %s → %s failed: %v", cam.ID, partFile, finalFile, err) } else { log.Printf("recorder: %s — segment %s complete (%v, %d bytes)", cam.ID, timeStr, elapsed.Round(time.Second), fileSize(finalFile)) } // Immediate loop — next segment starts with current timestamp. } } // fileSize returns the size of a file in bytes, or 0 on error. func fileSize(path string) int64 { info, err := os.Stat(path) if err != nil { return 0 } return info.Size() }