package service import ( "bufio" "context" "fmt" "io" "log/slog" "os/exec" "silk-server-go/internal/model" "gorm.io/gorm" ) // TranscodeService ffmpeg 实时转码服务 type TranscodeService struct { db *gorm.DB s3 *S3Service sem chan struct{} // 信号量,限制并发转码数为 4 } // NewTranscodeService 创建转码服务,信号量默认容量 4 func NewTranscodeService(db *gorm.DB, s3 *S3Service) *TranscodeService { return &TranscodeService{ db: db, s3: s3, sem: make(chan struct{}, 4), } } // StreamLive 通过 ffmpeg 实时转码直播流(H264 re-encode for browser compatibility) func (s *TranscodeService) StreamLive(sourceURL string, writer io.Writer) error { s.sem <- struct{}{} defer func() { <-s.sem }() ctx, cancel := context.WithCancel(context.Background()) defer cancel() cmd := exec.CommandContext(ctx, "ffmpeg", "-i", sourceURL, "-c:v", "libx264", "-preset", "ultrafast", "-tune", "zerolatency", "-b:v", "1M", "-maxrate", "1.5M", "-bufsize", "1M", "-g", "30", "-an", "-f", "flv", "pipe:1", ) stdoutPipe, err := cmd.StdoutPipe() if err != nil { return fmt.Errorf("创建 stdout 管道失败: %w", err) } stderrPipe, err := cmd.StderrPipe() if err != nil { return fmt.Errorf("创建 stderr 管道失败: %w", err) } if err := cmd.Start(); err != nil { return fmt.Errorf("启动 ffmpeg 失败: %w", err) } defer func() { if cmd.ProcessState == nil || !cmd.ProcessState.Exited() { _ = cmd.Process.Kill() } }() go func() { scanner := bufio.NewScanner(stderrPipe) for scanner.Scan() { slog.Info("ffmpeg-live", "msg", scanner.Text()) } }() if _, err := io.Copy(writer, stdoutPipe); err != nil { cancel() return fmt.Errorf("流式传输失败: %w", err) } if err := cmd.Wait(); err != nil { if ctx.Err() != nil { return ctx.Err() } return fmt.Errorf("ffmpeg 异常退出: %w", err) } return nil } // StreamClip 通过 ffmpeg 实时转码视频片段并流式写入 writer // 流程:查库获取 S3 信息 → 生成 presigned URL → ffmpeg -c copy 转封装为 fragmented MP4 → 管道输出 func (s *TranscodeService) StreamClip(clipId string, writer io.Writer) error { // 1. 查库获取 clip 的 s3Bucket + s3Key var clip model.VideoClip if err := s.db.Where("id = ?", clipId).First(&clip).Error; err != nil { return fmt.Errorf("视频片段不存在: %w", err) } if clip.S3Bucket == nil || clip.S3Key == nil || *clip.S3Bucket == "" || *clip.S3Key == "" { return fmt.Errorf("该片段没有 S3 对象") } // 2. 生成 S3 presigned URL presignedURL, err := s.s3.GetPresignedURL(*clip.S3Bucket, *clip.S3Key) if err != nil { return fmt.Errorf("生成 presigned URL 失败: %w", err) } // 3. 获取信号量(并发限制 4) s.sem <- struct{}{} defer func() { <-s.sem }() // 4. 启动 ffmpeg: ffmpeg -i -c copy -movflags frag_keyframe+empty_moov -f mp4 pipe:1 ctx, cancel := context.WithCancel(context.Background()) defer cancel() cmd := exec.CommandContext(ctx, "ffmpeg", "-i", presignedURL, "-c", "copy", "-movflags", "frag_keyframe+empty_moov", "-f", "mp4", "pipe:1", ) stdoutPipe, err := cmd.StdoutPipe() if err != nil { return fmt.Errorf("创建 stdout 管道失败: %w", err) } stderrPipe, err := cmd.StderrPipe() if err != nil { return fmt.Errorf("创建 stderr 管道失败: %w", err) } if err := cmd.Start(); err != nil { return fmt.Errorf("启动 ffmpeg 失败: %w", err) } // 确保 ffmpeg 进程被清理(如果还在运行则 kill) defer func() { if cmd.ProcessState == nil || !cmd.ProcessState.Exited() { _ = cmd.Process.Kill() } }() // 5. 记录 ffmpeg stderr 日志 go func() { scanner := bufio.NewScanner(stderrPipe) for scanner.Scan() { slog.Info("ffmpeg", "clipId", clipId, "msg", scanner.Text()) } }() // 6. 流式传输:io.Copy(writer, ffmpeg.Stdout) if _, err := io.Copy(writer, stdoutPipe); err != nil { cancel() // 取消 context,终止 ffmpeg return fmt.Errorf("流式传输失败: %w", err) } // 7. 等待 ffmpeg 结束,检查退出码 if err := cmd.Wait(); err != nil { if ctx.Err() != nil { return ctx.Err() } return fmt.Errorf("ffmpeg 异常退出: %w", err) } return nil }