diff --git a/.golangci.yaml b/.golangci.yaml index ea39029926..3a6223134d 100644 --- a/.golangci.yaml +++ b/.golangci.yaml @@ -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 diff --git a/internal/pkg/mediaconvert/command.go b/internal/pkg/mediaconvert/command.go index 88dbb7bb4a..3b02624031 100644 --- a/internal/pkg/mediaconvert/command.go +++ b/internal/pkg/mediaconvert/command.go @@ -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", diff --git a/internal/pkg/mediaconvert/command_test.go b/internal/pkg/mediaconvert/command_test.go index f8b36a7750..0b83a79ab6 100644 --- a/internal/pkg/mediaconvert/command_test.go +++ b/internal/pkg/mediaconvert/command_test.go @@ -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) +} diff --git a/internal/pkg/mediaconvert/transcoder.go b/internal/pkg/mediaconvert/transcoder.go index 6f476eee8d..7d8c0316d4 100644 --- a/internal/pkg/mediaconvert/transcoder.go +++ b/internal/pkg/mediaconvert/transcoder.go @@ -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)) } diff --git a/internal/pkg/queue/worker.go b/internal/pkg/queue/worker.go index 43e72c822d..7c9dae02b4 100644 --- a/internal/pkg/queue/worker.go +++ b/internal/pkg/queue/worker.go @@ -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 diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index aad4594e16..baffcaa82e 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -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)) }