package process import ( "bufio" "fmt" "io" "log" "os/exec" "strings" "sync" "time" "tuner/models" ) // FFmpegManager manages an FFmpeg child process for IPTV relay. type FFmpegManager struct { mu sync.Mutex cmd *exec.Cmd startTime time.Time lastError string streamURL string health models.StreamHealth } // NewFFmpegManager creates a new FFmpeg process manager. func NewFFmpegManager() *FFmpegManager { return &FFmpegManager{} } // Start spawns an FFmpeg process to pull from streamURL and push to MediaMTX. // It kills any existing FFmpeg process first. func (f *FFmpegManager) Start(streamURL string) error { f.mu.Lock() defer f.mu.Unlock() // Kill existing process if running f.stopLocked() f.cmd = exec.Command("ffmpeg", "-loglevel", "warning", "-progress", "pipe:1", "-reconnect", "1", "-reconnect_streamed", "1", "-reconnect_delay_max", "5", "-i", streamURL, "-c", "copy", "-f", "flv", "rtmp://localhost:1935/live/stream", ) f.streamURL = streamURL f.lastError = "" f.health = models.StreamHealth{} // Capture stdout for progress stats, stderr for errors stdout, err := f.cmd.StdoutPipe() if err != nil { return fmt.Errorf("stdout pipe: %w", err) } stderr, err := f.cmd.StderrPipe() if err != nil { return fmt.Errorf("stderr pipe: %w", err) } if err := f.cmd.Start(); err != nil { f.lastError = err.Error() return fmt.Errorf("start ffmpeg: %w", err) } f.startTime = time.Now() log.Printf("[ffmpeg] started with pid %d, source: %s", f.cmd.Process.Pid, streamURL) // Pass cmd reference so goroutines can detect process replacement cmd := f.cmd go f.parseProgress(stdout) go f.logStderr(stderr) go f.monitor(cmd) return nil } // parseProgress reads FFmpeg -progress output and updates health stats. func (f *FFmpegManager) parseProgress(r io.Reader) { scanner := bufio.NewScanner(r) for scanner.Scan() { line := scanner.Text() k, v, ok := strings.Cut(line, "=") if !ok { continue } f.mu.Lock() switch k { case "speed": f.health.Speed = v case "fps": f.health.FPS = v case "bitrate": f.health.Bitrate = v case "drop_frames": f.health.DropFrames = v case "dup_frames": f.health.DupFrames = v } f.mu.Unlock() } } func (f *FFmpegManager) logStderr(r io.Reader) { scanner := bufio.NewScanner(r) for scanner.Scan() { log.Printf("[ffmpeg] %s", scanner.Text()) } } func (f *FFmpegManager) monitor(cmd *exec.Cmd) { err := cmd.Wait() f.mu.Lock() defer f.mu.Unlock() // If the process was replaced by a new Start() call, don't touch state if f.cmd != cmd { return } if err != nil { f.lastError = err.Error() log.Printf("[ffmpeg] exited: %v", err) } else { log.Println("[ffmpeg] exited (status 0)") } f.cmd = nil } // Stop kills the FFmpeg process. func (f *FFmpegManager) Stop() error { f.mu.Lock() defer f.mu.Unlock() return f.stopLocked() } func (f *FFmpegManager) stopLocked() error { if f.cmd == nil || f.cmd.Process == nil { return nil } if err := f.cmd.Process.Kill(); err != nil { return fmt.Errorf("kill ffmpeg: %w", err) } // Wait for the process to fully exit to avoid zombies _ = f.cmd.Wait() f.cmd = nil log.Println("[ffmpeg] killed") return nil } // Health returns the latest stream health stats from FFmpeg progress. func (f *FFmpegManager) Health() *models.StreamHealth { f.mu.Lock() defer f.mu.Unlock() if f.cmd == nil || f.health.Speed == "" { return nil } h := f.health // copy return &h } // Status returns the current FFmpeg process status. func (f *FFmpegManager) Status() models.ProcessInfo { f.mu.Lock() defer f.mu.Unlock() info := models.ProcessInfo{ Error: f.lastError, } if f.cmd != nil && f.cmd.Process != nil { info.Running = true info.PID = f.cmd.Process.Pid info.Uptime = time.Since(f.startTime).Truncate(time.Second).String() } return info }