NextNVR/recorder.go

191 lines
5.7 KiB
Go

// NextNVR v0.2.1 — FFmpeg stream recorder
// Launches per-segment FFmpeg processes with stream-copy mode.
// Each recording runs for exactly 5 minutes (-t 300).
// In-progress files use .part.mp4 suffix; atomically renamed on completion.
// This prevents partial/corrupt files and provides crash recovery.
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 with start-of-segment timestamp.
now := time.Now()
// Round down to the nearest 5-minute boundary for clean naming.
rounded := now.Truncate(5 * time.Minute)
timeStr := rounded.Format("2006-01-02-15-04")
partFile := filepath.Join(outDir, timeStr+".part.mp4")
finalFile := filepath.Join(outDir, timeStr+".mp4")
// Build FFmpeg command with hard 5-minute timeout.
cmd := exec.Command("ffmpeg",
"-hide_banner", "-loglevel", "error",
"-rtsp_transport", "tcp",
"-i", rtspURL,
"-c", "copy", // stream-copy: 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()
}