fix: ensure correct ffmpeg commands

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>
This commit is contained in:
Alexander Onnikov
2025-06-24 16:44:22 +07:00
parent 05d9bb9cb9
commit 960eceb4af
6 changed files with 60 additions and 26 deletions
+2 -2
View File
@@ -61,8 +61,8 @@ linters-settings:
dupl:
threshold: 150
funlen:
lines: 200
statements: 100
lines: 240
statements: 120
goconst:
min-len: 2
min-occurrences: 2
+14 -1
View File
@@ -84,6 +84,8 @@ func buildCommonCommand(opts *Options) []string {
var result = []string{
"-y", // Overwrite output files without asking.
"-v", string(opts.LogLevel),
"-err_detect", "ignore_err",
"-fflags", "+discardcorrupt",
"-threads", fmt.Sprint(opts.Threads),
"-i", opts.Input,
}
@@ -113,6 +115,8 @@ func BuildAudioCommand(opts *Options) []string {
func BuildRawVideoCommand(opts *Options) []string {
if opts.Transcode {
return append(buildCommonCommand(opts),
"-map", "0:v:0",
"-map", "0:a?",
"-c:a", "aac",
"-c:v", "libx264",
"-preset", "veryfast",
@@ -149,6 +153,14 @@ func BuildThumbnailCommand(opts *Options) []string {
// BuildScalingVideoCommand returns flags for ffmpeg for video scaling
func BuildScalingVideoCommand(opts *Options) []string {
if len(opts.ScalingLevels) == 0 {
return []string{}
}
if len(opts.ScalingLevels) == 1 && opts.ScalingLevels[0] == opts.Level {
return []string{}
}
var result = buildCommonCommand(opts)
for _, level := range opts.ScalingLevels {
@@ -157,7 +169,8 @@ func BuildScalingVideoCommand(opts *Options) []string {
}
result = append(result,
"-map", "0:v",
"-map", "0:v:0",
"-map", "0:a?",
"-vf", "scale=-2:"+level[:len(level)-1],
"-c:a", "aac",
"-c:v", "libx264",
+18 -4
View File
@@ -32,7 +32,7 @@ func Test_BuildVideoCommand_Scaling(t *testing.T) {
ScalingLevels: []string{"720p", "480p"},
})
const expected = `-y -v debug -threads 4 -i pipe:0 -map 0:v -vf scale=-2:720 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_720p.ts test/1/1_720p_master.m3u8 -map 0:v -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
const expected = `-y -v debug -err_detect ignore_err -fflags +discardcorrupt -threads 4 -i pipe:0 -map 0:v:0 -map 0:a? -vf scale=-2:720 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_720p.ts test/1/1_720p_master.m3u8 -map 0:v:0 -map 0:a? -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
require.Contains(t, expected, strings.Join(scaleCommand, " "))
}
@@ -48,7 +48,7 @@ func Test_BuildVideoCommand_Scaling_NoRaw(t *testing.T) {
ScalingLevels: []string{"720p", "480p"},
})
const expected = `-y -v debug -threads 4 -i pipe:0 -map 0:v -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
const expected = `-y -v debug -err_detect ignore_err -fflags +discardcorrupt -threads 4 -i pipe:0 -map 0:v:0 -map 0:a? -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
require.Contains(t, expected, strings.Join(scaleCommand, " "))
}
@@ -64,7 +64,7 @@ func Test_BuildVideoCommand_Raw_NoTranscode(t *testing.T) {
Transcode: false,
})
const expected = `"-y -v debug -threads 4 -i pipe:0 -c:a copy -c:v copy -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
const expected = `"-y -v debug -err_detect ignore_err -fflags +discardcorrupt -threads 4 -i pipe:0 -c:a copy -c:v copy -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
require.Contains(t, expected, strings.Join(rawCommand, " "))
}
@@ -80,7 +80,21 @@ func Test_BuildVideoCommand_Raw_Transcode(t *testing.T) {
Transcode: true,
})
const expected = `-y -v debug -threads 4 -i pipe:0 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
const expected = `-y -v debug -err_detect ignore_err -fflags +discardcorrupt -threads 4 -i pipe:0 -map 0:v:0 -map 0:a? -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
require.Contains(t, expected, strings.Join(rawCommand, " "))
}
func Test_BuildVideoCommand_Scaling_Small(t *testing.T) {
var scaleCommand = mediaconvert.BuildScalingVideoCommand(&mediaconvert.Options{
OutputDir: "test",
Input: "pipe:0",
UploadID: "1",
Threads: 4,
LogLevel: mediaconvert.LogLevelDebug,
Level: "360p",
ScalingLevels: []string{"360p"},
})
require.Empty(t, scaleCommand)
}
+12 -7
View File
@@ -47,8 +47,8 @@ type Transcoder struct {
// Command represents a ffmpeg command
type Command struct {
cmd *exec.Cmd
stdoutBuf bytes.Buffer
stderrBuf bytes.Buffer
stdoutBuf *bytes.Buffer
stderrBuf *bytes.Buffer
}
// NewTranscoder creates a new instance of task transcoder
@@ -183,6 +183,11 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er
var cmds []Command
for _, args := range argsSlice {
if len(args) == 0 {
logger.Debug("skip empty command")
continue
}
cmd, cmdErr := newFfmpegCommand(ctx, nil, args)
if cmdErr != nil {
logger.Error("can not create a new command", zap.Error(cmdErr), zap.Strings("args", args))
@@ -192,12 +197,12 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er
var command = Command{
cmd: cmd,
stdoutBuf: bytes.Buffer{},
stderrBuf: bytes.Buffer{},
stdoutBuf: &bytes.Buffer{},
stderrBuf: &bytes.Buffer{},
}
cmd.Stdout = io.MultiWriter(os.Stdout, &command.stdoutBuf)
cmd.Stderr = io.MultiWriter(os.Stderr, &command.stderrBuf)
cmd.Stdout = io.MultiWriter(os.Stdout, command.stdoutBuf)
cmd.Stderr = io.MultiWriter(os.Stderr, command.stderrBuf)
cmds = append(cmds, command)
if startErr := cmd.Start(); startErr != nil {
@@ -215,7 +220,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er
continue
}
logger.Error("can not wait for command end ", zap.Error(cmdErr))
logger.Error("can not wait for command end", zap.Error(cmdErr), zap.String("cmd", cmd.cmd.String()))
if _, writeErr := os.Stdout.Write(cmd.stdoutBuf.Bytes()); writeErr != nil {
logger.Error("can not write stdout ", zap.Error(writeErr))
}
+1 -1
View File
@@ -85,7 +85,7 @@ func (w *Worker) fetchAndProcessMessage(ctx context.Context) error {
err = w.processMessage(ctx, msg, logger)
if err != nil {
w.logger.Error("failed to process message", zap.Error(err))
return err
return fmt.Errorf("process message: %w", err)
}
return nil
+13 -11
View File
@@ -297,6 +297,13 @@ func (u *uploaderImpl) uploadAndDelete(f string) {
return
}
// Check if the file has already been uploaded
_, ok := u.sentFiles.Load(f)
if ok && !u.shouldDeleteOnStop(f) {
logger.Debug("file already uploaded")
return
}
if err := waitFileExists(f); err != nil {
if os.IsNotExist(err) {
logger.Debug("file does not exist", zap.Error(err))
@@ -306,18 +313,11 @@ func (u *uploaderImpl) uploadAndDelete(f string) {
return
}
// Check if the file has already been uploaded
var _, ok = u.sentFiles.Load(f)
if ok && !u.shouldDeleteOnStop(f) {
logger.Debug("file already uploaded")
return
}
for attempt := range u.options.RetryCount {
logger = logger.With(zap.Int("attempt", attempt))
var ctx, cancel = context.WithTimeout(u.uploadCtx, u.options.Timeout)
var err = u.storage.PutFile(ctx, f)
cancel()
var putCtx, putCancel = context.WithTimeout(u.uploadCtx, u.options.Timeout)
var err = u.storage.PutFile(putCtx, f)
putCancel()
if err != nil {
logger.Error("attempt failed", zap.Error(err))
@@ -328,7 +328,9 @@ func (u *uploaderImpl) uploadAndDelete(f string) {
// Update the file's parent if SourceFile is set
if u.options.Source != "" {
err = u.storage.SetParent(ctx, f, u.options.Source)
var setParentCtx, setParentCancel = context.WithTimeout(u.uploadCtx, u.options.Timeout)
err = u.storage.SetParent(setParentCtx, f, u.options.Source)
setParentCancel()
if err != nil {
logger.Error("can not set blob parent", zap.Error(err), zap.String("filename", f), zap.String("source", u.options.Source))
}