NextNVR/recorder.go

141 lines
3.7 KiB
Go

// NextNVR v0.2.0 — FFmpeg stream recorder
// Launches per-camera FFmpeg processes with stream-copy mode and auto-reconnect.
package main
import (
"fmt"
"log"
"os"
"os/exec"
"path/filepath"
"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 launches an FFmpeg process for a single camera.
// Uses stream-copy mode (-c copy) for zero transcoding overhead.
// Segments output into 5-minute MP4 files.
// Auto-reconnects on stream failure with a 5-second delay.
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()
for {
select {
case <-proc.done:
return
default:
}
// Create output directory: /mnt/recordings/{cam_id}/{YYYY-MM-DD}/
now := time.Now()
outDir := filepath.Join(rm.config.RecordingsPath, cam.ID, now.Format("2006-01-02"))
if err := os.MkdirAll(outDir, 0755); err != nil {
log.Printf("recorder: %s — mkdir failed: %v", cam.ID, err)
time.Sleep(5 * time.Second)
continue
}
// Output filename: {HH-MM-SS}.mp4
outFile := filepath.Join(outDir, now.Format("15-04-05")+".mp4")
cmd := exec.Command("ffmpeg",
"-hide_banner", "-loglevel", "error",
"-rtsp_transport", "tcp", // TCP is more reliable than UDP for RTSP
"-i", rtspURL,
"-c", "copy", // stream-copy: zero transcoding
"-f", "segment",
"-segment_time", "300", // 5-minute segments
"-segment_format", "mp4",
"-reset_timestamps", "1",
"-strftime", "1",
outFile,
)
cmd.Stdout = os.Stdout
cmd.Stderr = os.Stderr
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 — FFmpeg running (PID %d)", cam.ID, cmd.Process.Pid)
err := cmd.Wait()
proc.cmd = nil
if err != nil {
log.Printf("recorder: %s — FFmpeg exited: %v — reconnecting in 5s", cam.ID, err)
} else {
log.Printf("recorder: %s — FFmpeg exited cleanly — reconnecting in 5s", cam.ID)
}
time.Sleep(5 * time.Second)
}
}