192 lines
5.8 KiB
Go
192 lines
5.8 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.
|
|
// All timestamps are in local time (BST) — no UTC conversion.
|
|
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()
|
|
}
|