From c83ff3bf6837eb2b3d4a3530a82fb88c2595ca5e Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Fri, 20 Jun 2025 16:25:10 +0700 Subject: [PATCH 1/3] fix: properly detect created hls segments Signed-off-by: Alexander Onnikov --- Dockerfile | 2 +- go.mod | 2 +- go.sum | 4 +- internal/pkg/mediaconvert/command.go | 34 +++++++++++--- internal/pkg/mediaconvert/command_test.go | 8 ++-- internal/pkg/mediaconvert/scheduler.go | 2 + internal/pkg/mediaconvert/transcoder.go | 56 ++++++++++++++++------- internal/pkg/storage/datalake.go | 56 +++++++++++++---------- internal/pkg/uploader/uploader.go | 25 ++++++++-- 9 files changed, 131 insertions(+), 58 deletions(-) diff --git a/Dockerfile b/Dockerfile index 121a1b5efb..1e7a5cd10b 100644 --- a/Dockerfile +++ b/Dockerfile @@ -11,7 +11,7 @@ # See the License for the specific language governing permissions and # limitations under the License. -FROM golang:1.24.1 AS builder +FROM golang:1.24.4 AS builder ENV GO111MODULE=on ENV CGO_ENABLED=0 ENV GOBIN=/bin diff --git a/go.mod b/go.mod index 905e524afa..701291bda0 100644 --- a/go.mod +++ b/go.mod @@ -8,7 +8,7 @@ require ( github.com/aws/aws-sdk-go-v2/credentials v1.17.59 github.com/aws/aws-sdk-go-v2/service/s3 v1.77.0 github.com/getsentry/sentry-go v0.31.1 - github.com/golang-jwt/jwt/v5 v5.2.1 + github.com/golang-jwt/jwt/v5 v5.2.2 github.com/google/uuid v1.6.0 github.com/kelseyhightower/envconfig v1.4.0 github.com/pkg/errors v0.9.1 diff --git a/go.sum b/go.sum index c531b49b39..3ca0c0a4f2 100644 --- a/go.sum +++ b/go.sum @@ -45,8 +45,8 @@ github.com/getsentry/sentry-go v0.31.1 h1:ELVc0h7gwyhnXHDouXkhqTFSO5oslsRDk0++ey github.com/getsentry/sentry-go v0.31.1/go.mod h1:CYNcMMz73YigoHljQRG+qPF+eMq8gG72XcGN/p71BAY= github.com/go-errors/errors v1.4.2 h1:J6MZopCL4uSllY1OfXM374weqZFFItUbrImctkmUxIA= github.com/go-errors/errors v1.4.2/go.mod h1:sIVyrIiJhuEF+Pj9Ebtd6P/rEYROXFi3BopGUQ5a5Og= -github.com/golang-jwt/jwt/v5 v5.2.1 h1:OuVbFODueb089Lh128TAcimifWaLhJwVflnrgM17wHk= -github.com/golang-jwt/jwt/v5 v5.2.1/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= +github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8= +github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc= github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= diff --git a/internal/pkg/mediaconvert/command.go b/internal/pkg/mediaconvert/command.go index 047aa3a5a7..2efbf54cf7 100644 --- a/internal/pkg/mediaconvert/command.go +++ b/internal/pkg/mediaconvert/command.go @@ -19,7 +19,6 @@ import ( "context" "fmt" "io" - "os" "os/exec" "path/filepath" "strings" @@ -30,12 +29,27 @@ import ( "go.uber.org/zap" ) +type LogLevel string + +const ( + LogLevelQuiet LogLevel = "quiet" + LogLevelPanic LogLevel = "panic" + LogLevelFatal LogLevel = "fatal" + LogLevelError LogLevel = "error" + LogLevelWarning LogLevel = "warning" + LogLevelInfo LogLevel = "info" + LogLevelVerbose LogLevel = "verbose" + LogLevelDebug LogLevel = "debug" + LogLevelTrace LogLevel = "trace" +) + // Options represents configuration for the ffmpeg command type Options struct { Input string OutputDir string ScalingLevels []string Level string + LogLevel LogLevel Transcode bool Threads int UploadID string @@ -51,8 +65,6 @@ func newFfmpegCommand(ctx context.Context, in io.Reader, args []string) (*exec.C logger.Debug("prepared command: ", zap.Strings("args", args)) var result = exec.CommandContext(ctx, "ffmpeg", args...) - result.Stderr = os.Stdout - result.Stdout = os.Stdout result.Stdin = in return result, nil @@ -60,6 +72,8 @@ func newFfmpegCommand(ctx context.Context, in io.Reader, args []string) (*exec.C func buildCommonCommand(opts *Options) []string { var result = []string{ + "-y", // Overwrite output files without asking. + "-v", string(opts.LogLevel), "-threads", fmt.Sprint(opts.Threads), "-i", opts.Input, } @@ -96,7 +110,7 @@ func BuildRawVideoCommand(opts *Options) []string { "-g", "60", "-f", "hls", "-hls_time", "5", - "-hls_flags", "split_by_time", + "-hls_flags", "split_by_time+temp_file", "-hls_list_size", "0", "-hls_segment_filename", filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", opts.Level)), filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, opts.Level))) @@ -107,7 +121,7 @@ func BuildRawVideoCommand(opts *Options) []string { "-c:v", "copy", // Copy video stream "-f", "hls", "-hls_time", "5", - "-hls_flags", "split_by_time", + "-hls_flags", "split_by_time+temp_file", "-hls_list_size", "0", "-hls_segment_filename", filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", opts.Level)), filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, opts.Level))) @@ -142,7 +156,15 @@ func BuildScalingVideoCommand(opts *Options) []string { "-g", "60", "-f", "hls", "-hls_time", "5", - "-hls_flags", "split_by_time", + // Use HLS flags + // - split_by_time + // Allow segments to start on frames other than key frames. + // This improves behavior on some players when the time between key frames is inconsistent, + // but may make things worse on others, and can cause some oddities during seeking. + // This flag should be used with the hls_time option. + // - temp_file + // Write segment data to filename.tmp and rename to filename only once the segment is complete. + "-hls_flags", "split_by_time+temp_file", "-hls_list_size", "0", "-hls_segment_filename", filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", level)), filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, level))) diff --git a/internal/pkg/mediaconvert/command_test.go b/internal/pkg/mediaconvert/command_test.go index b6161e7ea7..26beed3272 100644 --- a/internal/pkg/mediaconvert/command_test.go +++ b/internal/pkg/mediaconvert/command_test.go @@ -31,7 +31,7 @@ func Test_BuildVideoCommand_Scaling(t *testing.T) { ScalingLevels: []string{"720p", "480p"}, }) - const expected = `-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 -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 -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8` + const expected = `-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 -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, " ")) } @@ -46,7 +46,7 @@ func Test_BuildVideoCommand_Scaling_NoRaw(t *testing.T) { ScalingLevels: []string{"720p", "480p"}, }) - const expected = `-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 -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8` + const expected = `-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` require.Contains(t, expected, strings.Join(scaleCommand, " ")) } @@ -61,7 +61,7 @@ func Test_BuildVideoCommand_Raw_NoTranscode(t *testing.T) { Transcode: false, }) - const expected = `"-threads 4 -i pipe:0 -c:a copy -c:v copy -f hls -hls_time 5 -hls_flags split_by_time -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8` + const expected = `"-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, " ")) } @@ -76,7 +76,7 @@ func Test_BuildVideoCommand_Raw_Transcode(t *testing.T) { Transcode: true, }) - const expected = `-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 -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8` + const expected = `-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` require.Contains(t, expected, strings.Join(rawCommand, " ")) } diff --git a/internal/pkg/mediaconvert/scheduler.go b/internal/pkg/mediaconvert/scheduler.go index 1d71d9113f..5f33537360 100644 --- a/internal/pkg/mediaconvert/scheduler.go +++ b/internal/pkg/mediaconvert/scheduler.go @@ -299,6 +299,8 @@ func IsSupportedMediaType(mediaType string) bool { return true case "video/webm": return true + case "video/quicktime": + return true default: return false } diff --git a/internal/pkg/mediaconvert/transcoder.go b/internal/pkg/mediaconvert/transcoder.go index 8b04ca677d..fca9d02353 100644 --- a/internal/pkg/mediaconvert/transcoder.go +++ b/internal/pkg/mediaconvert/transcoder.go @@ -16,8 +16,10 @@ package mediaconvert import ( + "bytes" "context" "fmt" + "io" "os" "os/exec" "path/filepath" @@ -42,6 +44,12 @@ type Transcoder struct { logger *zap.Logger } +type Command struct { + cmd *exec.Cmd + stdoutBuf bytes.Buffer + stderrBuf bytes.Buffer +} + // NewTranscoder creates a new instance of task transcoder func NewTranscoder(ctx context.Context, cfg *config.Config) *Transcoder { var p = &Transcoder{ @@ -54,7 +62,7 @@ func NewTranscoder(ctx context.Context, cfg *config.Config) *Transcoder { } // Transcode handles one transcoding task -func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, error) { +func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, error) { var logger = p.logger.With(zap.String("task-id", task.ID)) logger.Debug("start") @@ -64,7 +72,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err var tokenString, err = token.NewToken(p.cfg.ServerSecret, task.Workspace, "stream", "datalake") if err != nil { logger.Error("can not create token", zap.Error(err)) - return TaskResult{}, errors.Wrapf(err, "can not create token") + return nil, errors.Wrapf(err, "can not create token") } logger.Debug("phase 2: preparing fs") @@ -73,7 +81,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err err = os.MkdirAll(destinationFolder, os.ModePerm) if err != nil { logger.Error("can not create temporary folder", zap.Error(err)) - return TaskResult{}, errors.Wrapf(err, "can not create temporary folder") + return nil, errors.Wrapf(err, "can not create temporary folder") } defer func() { @@ -87,38 +95,38 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err remoteStorage, err := storage.NewStorageByURL(ctx, p.cfg.Endpoint(), p.cfg.EndpointURL.Scheme, tokenString, task.Workspace) if err != nil { logger.Error("can not create storage by url", zap.Error(err), zap.String("url", p.cfg.EndpointURL.String())) - return TaskResult{}, errors.Wrapf(err, "can not create storage by url") + return nil, errors.Wrapf(err, "can not create storage by url") } stat, err := remoteStorage.StatFile(ctx, task.Source) if err != nil { logger.Error("can not stat file", zap.Error(err), zap.String("filepath", task.Source)) - return TaskResult{}, errors.Wrapf(err, "can not stat file") + return nil, errors.Wrapf(err, "can not stat file") } if !IsSupportedMediaType(stat.Type) { logger.Info("unsupported media type", zap.String("type", stat.Type)) - return TaskResult{}, fmt.Errorf("unsupported media type: %s", stat.Type) + return nil, fmt.Errorf("unsupported media type: %s", stat.Type) } sourceFilePath := filepath.Join(destinationFolder, filename) if err = remoteStorage.GetFile(ctx, task.Source, sourceFilePath); err != nil { logger.Error("can not download source file", zap.Error(err), zap.String("filepath", task.Source)) // TODO: reschedule - return TaskResult{}, errors.Wrapf(err, "can not download source file") + return nil, errors.Wrapf(err, "can not download source file") } logger.Debug("phase 4: prepare to transcode") probe, err := ffprobe.ProbeURL(ctx, sourceFilePath) if err != nil { logger.Error("can not get ffprobe", zap.Error(err), zap.String("filepath", sourceFilePath)) - return TaskResult{}, errors.Wrapf(err, "can not get ffprobe") + return nil, errors.Wrapf(err, "can not get ffprobe") } videoStream := probe.FirstVideoStream() if videoStream == nil { logger.Error("no video stream found", zap.String("filepath", sourceFilePath)) - return TaskResult{}, errors.Wrapf(err, "no video stream found") + return nil, errors.Wrapf(err, "no video stream found") } logger.Debug("video stream found", zap.String("codec", videoStream.CodecName), zap.Int("width", videoStream.Width), zap.Int("height", videoStream.Height)) @@ -136,6 +144,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err Input: sourceFilePath, OutputDir: p.cfg.OutputDir, Level: level, + LogLevel: LogLevel(p.cfg.LogLevel), Transcode: !IsHLSSupportedVideoCodec(codec), ScalingLevels: append(sublevels, level), UploadID: task.ID, @@ -157,7 +166,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err err = manifest.GenerateHLSPlaylist(opts.ScalingLevels, p.cfg.OutputDir, opts.UploadID) if err != nil { logger.Error("can not generate hls playlist", zap.String("out", p.cfg.OutputDir), zap.String("uploadID", opts.UploadID)) - return TaskResult{}, errors.Wrapf(err, "can not generate hls playlist") + return nil, errors.Wrapf(err, "can not generate hls playlist") } go uploader.Start() @@ -169,29 +178,42 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err BuildRawVideoCommand(&opts), BuildScalingVideoCommand(&opts), } - var cmds []*exec.Cmd + var cmds []Command for _, args := range argsSlice { cmd, cmdErr := newFfmpegCommand(ctx, nil, args) if cmdErr != nil { logger.Error("can not create a new command", zap.Error(cmdErr), zap.Strings("args", args)) go uploader.Cancel() - return TaskResult{}, errors.Wrapf(err, "can not create a new command") + return nil, errors.Wrapf(err, "can not create a new command") } - cmds = append(cmds, cmd) + + var command = Command{ + cmd: cmd, + stdoutBuf: bytes.Buffer{}, + stderrBuf: bytes.Buffer{}, + } + + cmd.Stdout = io.MultiWriter(os.Stdout, &command.stdoutBuf) + cmd.Stderr = io.MultiWriter(os.Stderr, &command.stderrBuf) + + cmds = append(cmds, command) if err = cmd.Start(); err != nil { logger.Error("can not start a command", zap.Error(err), zap.Strings("args", args)) go uploader.Cancel() - return TaskResult{}, errors.Wrapf(err, "can not start a command") + return nil, errors.Wrapf(err, "can not start a command") } } logger.Debug("phase 7: wait for the result") + for _, cmd := range cmds { - if err = cmd.Wait(); err != nil { + if err = cmd.cmd.Wait(); err != nil { logger.Error("can not wait for command end ", zap.Error(err)) + os.Stdout.Write(cmd.stdoutBuf.Bytes()) + os.Stderr.Write(cmd.stderrBuf.Bytes()) go uploader.Cancel() - return TaskResult{}, errors.Wrapf(err, "can not wait for command end") + return nil, errors.Wrapf(err, "can not wait for command end") } } @@ -226,5 +248,5 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err } } - return result, nil + return &result, nil } diff --git a/internal/pkg/storage/datalake.go b/internal/pkg/storage/datalake.go index 8a5284197f..aaf18ac273 100644 --- a/internal/pkg/storage/datalake.go +++ b/internal/pkg/storage/datalake.go @@ -21,8 +21,10 @@ import ( "io" "mime/multipart" "net/textproto" + "net/url" "os" "path/filepath" + "strconv" "strings" "time" @@ -170,6 +172,11 @@ func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error return errors.Wrapf(err, "delete failed") } + if err := okResponse(res); err != nil { + logRequestError(logger, err, "bad status code", res) + return err + } + logger.Debug("deleted") return nil @@ -205,14 +212,11 @@ func (d *DatalakeStorage) PatchMeta(ctx context.Context, filename string, md *Me return err } - if resp.StatusCode() != fasthttp.StatusOK { - var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode()) - logger.Debug("bad status code", zap.Error(err)) + if err := okResponse(resp); err != nil { + logRequestError(logger, err, "bad status code", resp) return err } - fmt.Println(string(resp.Body())) - return nil } @@ -237,9 +241,8 @@ func (d *DatalakeStorage) GetMeta(ctx context.Context, filename string) (*Metada return nil, err } - if resp.StatusCode() != fasthttp.StatusOK { - var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode()) - logger.Debug("bad status code", zap.Error(err)) + if err := okResponse(resp); err != nil { + logRequestError(logger, err, "bad status code", resp) return nil, err } @@ -270,10 +273,8 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str return err } - // Check the response status code - if resp.StatusCode() != fasthttp.StatusOK { - var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode()) - logger.Debug("bad status code", zap.Error(err)) + if err := okResponse(resp); err != nil { + logRequestError(logger, err, "bad status code", resp) return err } @@ -291,7 +292,13 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str return err } - logger.Debug("file downloaded successfully") + stat, err := os.Stat(destination) + if err != nil { + logger.Error("can't stat the file", zap.Error(err)) + return err + } + + logger.Info("file downloaded successfully", zap.Int64("size", stat.Size())) return nil } @@ -316,9 +323,7 @@ func (d *DatalakeStorage) StatFile(ctx context.Context, filename string) (*BlobI return nil, err } - // Check the response status code - if resp.StatusCode() != fasthttp.StatusOK { - var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode()) + if err := okResponse(resp); err != nil { logRequestError(logger, err, "bad status code", resp) return nil, err } @@ -346,7 +351,7 @@ func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string req.SetRequestURI(d.baseURL + "/blob/" + d.workspace + "/" + objectKey + "/parent") req.Header.SetMethod(fasthttp.MethodPatch) req.Header.Add("Authorization", "Bearer "+d.token) - req.Header.Add("Content-Type", "application/json") + req.Header.SetContentType("application/json") body := map[string]any{ "parent": parentKey, @@ -365,15 +370,20 @@ func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string return err } - // Check the response status code - var statusOK = resp.StatusCode() >= 200 && resp.StatusCode() < 300 - if !statusOK { - var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode()) - logger.Debug("bad status code", zap.Error(err), zap.Int("status", resp.StatusCode()), zap.String("response", resp.String())) + if err := okResponse(resp); err != nil { + logRequestError(logger, err, "bad status code", resp) return err } - logger.Debug("finished") + return nil +} + +func okResponse(res *fasthttp.Response) error { + var statusOK = res.StatusCode() >= 200 && res.StatusCode() < 300 + + if !statusOK { + return fmt.Errorf("unexpected status code: %d", res.StatusCode()) + } return nil } diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index 957d889f61..1120751029 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -35,6 +35,8 @@ import ( // See at https://man7.org/linux/man-pages/man7/inotify.7.html const inotifyCloseWrite uint32 = 0x8 // IN_CLOSE_WRITE const inotifyMovedTo uint32 = 0x80 // IN_MOVED_TO +const inotifyDelete uint32 = 0x200 // IN_DELETE +const inotifyMovedFrom uint32 = 0x40 // IN_MOVED_FROM // Uploader represents file uploader type Uploader interface { @@ -250,9 +252,12 @@ func (u *uploaderImpl) uploadAndDelete(f string) { } // Check if the file exists - _, err := os.Stat(f) - if err != nil { - logger.Debug("file does not exist") + if _, err := os.Stat(f); err != nil { + if os.IsNotExist(err) { + logger.Debug("file does not exist", zap.Error(err)) + } else { + logger.Error("failed to stat file", zap.Error(err)) + } return } @@ -300,6 +305,7 @@ func (u *uploaderImpl) uploadAndDelete(f string) { } } +// startWatch watches for changes in the directory and uploads created files func (u *uploaderImpl) startWatch(ready chan<- struct{}) { defer close(u.watcherDoneCh) @@ -317,7 +323,7 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) { } }() - if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite|inotifyMovedTo); err != nil { + if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite|inotifyMovedTo|inotifyDelete|inotifyMovedFrom); err != nil { logger.Error("can not start watching", zap.Error(err)) close(ready) return @@ -348,8 +354,19 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) { if strings.HasSuffix(event.Name, ".tmp") { continue } + if event.Mask&(inotifyDelete|inotifyMovedFrom) != 0 { + logger.Debug("file deleted or moved away", zap.String("event", event.Name), zap.Uint32("mask", event.Mask)) + continue + } + logger.Debug("received an event", zap.String("event", event.Name), zap.Uint32("mask", event.Mask)) + if _, err := os.Stat(event.Name); os.IsNotExist(err) { + logger.Warn("file does not exist", zap.String("file", event.Name)) + // wait a bit for file operations to complete + time.Sleep(100 * time.Millisecond) + } + u.filesCh <- event.Name case err, ok := <-watcher.Error: if !ok { From c910668a86ee59b9a265dd87c1ea5f6abe4b0a34 Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Mon, 23 Jun 2025 00:05:59 +0700 Subject: [PATCH 2/3] fix: wait until uploader finishes Signed-off-by: Alexander Onnikov --- internal/pkg/mediaconvert/transcoder.go | 3 +- internal/pkg/storage/datalake.go | 2 - internal/pkg/uploader/uploader.go | 126 +++++++++++++++--------- 3 files changed, 79 insertions(+), 52 deletions(-) diff --git a/internal/pkg/mediaconvert/transcoder.go b/internal/pkg/mediaconvert/transcoder.go index fca9d02353..46b8206083 100644 --- a/internal/pkg/mediaconvert/transcoder.go +++ b/internal/pkg/mediaconvert/transcoder.go @@ -85,6 +85,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er } defer func() { + logger.Debug("remove temporary folder") if err = os.RemoveAll(destinationFolder); err != nil { logger.Error("failed to cleanup temporary folder", zap.Error(err)) } @@ -218,7 +219,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er } logger.Debug("phase 8: schedule cleanup") - go uploader.Stop() + uploader.Stop() logger.Debug("phase 9: try to set metadata") diff --git a/internal/pkg/storage/datalake.go b/internal/pkg/storage/datalake.go index aaf18ac273..33d8b9e159 100644 --- a/internal/pkg/storage/datalake.go +++ b/internal/pkg/storage/datalake.go @@ -21,10 +21,8 @@ import ( "io" "mime/multipart" "net/textproto" - "net/url" "os" "path/filepath" - "strconv" "strings" "time" diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index 1120751029..e14782f0a7 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -34,9 +34,9 @@ import ( // See at https://man7.org/linux/man-pages/man7/inotify.7.html const inotifyCloseWrite uint32 = 0x8 // IN_CLOSE_WRITE +const inotifyMovedFrom uint32 = 0x40 // IN_MOVED_FROM const inotifyMovedTo uint32 = 0x80 // IN_MOVED_TO const inotifyDelete uint32 = 0x200 // IN_DELETE -const inotifyMovedFrom uint32 = 0x40 // IN_MOVED_FROM // Uploader represents file uploader type Uploader interface { @@ -104,50 +104,66 @@ func New(ctx context.Context, s storage.Storage, opts Options) Uploader { } func (u *uploaderImpl) Stop() { + u.logger.Info("stopping upload") u.stop(false) } func (u *uploaderImpl) Cancel() { + u.logger.Info("canceling upload") u.stop(true) } -func (u *uploaderImpl) scanInitialFiles() { - u.workerWaitGroup.Add(1) +func (u *uploaderImpl) scanFiles() { + logger := u.logger.With(zap.String("dir", u.options.Dir)) - go func() { - defer u.workerWaitGroup.Done() + logger.Info("scan files") + files, err := os.ReadDir(u.options.Dir) + if err != nil { + logger.Error("failed to read files", zap.Error(err)) + return + } - logger := u.logger.With(zap.String("dir", u.options.Dir)) - - logger.Info("initial file scan") - initFiles, err := os.ReadDir(u.options.Dir) - if err != nil { - logger.Error("failed to read initial files", zap.Error(err)) - return + count := 0 + for _, f := range files { + if f.IsDir() { + continue } - for _, f := range initFiles { - if f.IsDir() { - continue - } + var filePath = filepath.Join(u.options.Dir, f.Name()) - // Ignore source file - var filePath = filepath.Join(u.options.Dir, f.Name()) - if filePath == u.options.SourceFile { - continue - } - u.filesCh <- filePath + // Ignore source file + if filePath == u.options.SourceFile { + continue } - logger.Info("initial file scan complete", zap.Int("count", len(initFiles))) - }() + if _, uploaded := u.sentFiles.Load(filePath); uploaded { + logger.Debug("file already uploaded", zap.String("file", filePath)) + continue + } + + u.filesCh <- filePath + count++ + } + + logger.Info("scan complete", zap.Int("count", count)) } func (u *uploaderImpl) stop(rollback bool) { + // Stop watching for new files close(u.watcherStopCh) <-u.watcherDoneCh - u.logger.Debug("file watch stopped") + // Scan remaining files in the directory + u.scanFiles() + + // Close filesCh so no new files added + close(u.filesCh) + + // Wait for all workers to finish processing + u.workerWaitGroup.Wait() + u.logger.Debug("workers done") + + // Perform rollback if rollback { u.logger.Debug("starting rollback...") var i uint32 @@ -161,15 +177,25 @@ func (u *uploaderImpl) stop(rollback bool) { }) u.logger.Debug("rollback done") } - close(u.filesCh) - u.workerWaitGroup.Wait() - u.logger.Debug("workers done") u.uploadCancel() - _ = os.RemoveAll(u.options.Dir) + + remainingFiles, err := os.ReadDir(u.options.Dir) + if err != nil && !os.IsNotExist(err) { + u.logger.Error("failed to read dir", zap.Error(err)) + } + // log remaining files + if len(remainingFiles) > 0 { + files := make([]string, 0, len(remainingFiles)) + for _, entry := range remainingFiles { + files = append(files, entry.Name()) + } + u.logger.Info("remaining files", zap.Int("count", len(files)), zap.Any("files", files)) + } + u.sentFiles.Clear() - u.logger.Debug("finish done", zap.Bool("cancel", rollback)) + u.logger.Debug("stopped", zap.Bool("rollback", rollback)) } func (u *uploaderImpl) Start() { @@ -180,7 +206,7 @@ func (u *uploaderImpl) Start() { <-watcherReady - u.scanInitialFiles() + u.scanFiles() } func (u *uploaderImpl) startWorkers() { @@ -251,8 +277,7 @@ func (u *uploaderImpl) uploadAndDelete(f string) { return } - // Check if the file exists - if _, err := os.Stat(f); err != nil { + if err := waitFileExists(f); os.IsNotExist(err) { if os.IsNotExist(err) { logger.Debug("file does not exist", zap.Error(err)) } else { @@ -342,18 +367,12 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) { logger.Error("file channel was closed") return } - if !strings.Contains(event.Name, u.options.Dir) { - continue - } - if event.Name == u.options.Dir { - continue - } - if event.Name == u.options.SourceFile { - continue - } - if strings.HasSuffix(event.Name, ".tmp") { + if event.Name == u.options.Dir || + event.Name == u.options.SourceFile || + strings.HasSuffix(event.Name, ".tmp") { continue } + if event.Mask&(inotifyDelete|inotifyMovedFrom) != 0 { logger.Debug("file deleted or moved away", zap.String("event", event.Name), zap.Uint32("mask", event.Mask)) continue @@ -361,12 +380,6 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) { logger.Debug("received an event", zap.String("event", event.Name), zap.Uint32("mask", event.Mask)) - if _, err := os.Stat(event.Name); os.IsNotExist(err) { - logger.Warn("file does not exist", zap.String("file", event.Name)) - // wait a bit for file operations to complete - time.Sleep(100 * time.Millisecond) - } - u.filesCh <- event.Name case err, ok := <-watcher.Error: if !ok { @@ -376,3 +389,18 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) { } } } + +func waitFileExists(file string) error { + var err error + + for range 10 { + stat, err := os.Stat(file) + if err == nil && stat.Size() > 0 { + return nil + } + + time.Sleep(50 * time.Millisecond) + } + + return err +} From 0755e927f56fb89c19d624336841810fb0801f66 Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Mon, 23 Jun 2025 00:34:51 +0700 Subject: [PATCH 3/3] fix tests and lint issues Signed-off-by: Alexander Onnikov --- .golangci.yaml | 4 ++-- internal/pkg/mediaconvert/command.go | 24 ++++++++++++++++------- internal/pkg/mediaconvert/command_test.go | 12 ++++++++---- internal/pkg/mediaconvert/transcoder.go | 21 +++++++++++++------- internal/pkg/storage/datalake.go | 2 +- internal/pkg/uploader/uploader.go | 5 +++-- 6 files changed, 45 insertions(+), 23 deletions(-) diff --git a/.golangci.yaml b/.golangci.yaml index 5fd9dcf2b3..ea39029926 100644 --- a/.golangci.yaml +++ b/.golangci.yaml @@ -57,11 +57,11 @@ linters-settings: goimports: local-prefixes: github.com/networkservicemesh/sdk gocyclo: - min-complexity: 20 + min-complexity: 30 dupl: threshold: 150 funlen: - lines: 180 + lines: 200 statements: 100 goconst: min-len: 2 diff --git a/internal/pkg/mediaconvert/command.go b/internal/pkg/mediaconvert/command.go index 2efbf54cf7..88dbb7bb4a 100644 --- a/internal/pkg/mediaconvert/command.go +++ b/internal/pkg/mediaconvert/command.go @@ -29,18 +29,28 @@ import ( "go.uber.org/zap" ) +// LogLevel is ffmpeg log level type LogLevel string const ( - LogLevelQuiet LogLevel = "quiet" - LogLevelPanic LogLevel = "panic" - LogLevelFatal LogLevel = "fatal" - LogLevelError LogLevel = "error" + // LogLevelQuiet is quiet log level + LogLevelQuiet LogLevel = "quiet" + // LogLevelPanic is panic log level + LogLevelPanic LogLevel = "panic" + // LogLevelFatal is fatal log level + LogLevelFatal LogLevel = "fatal" + // LogLevelError is error log level + LogLevelError LogLevel = "error" + // LogLevelWarning is warning log level LogLevelWarning LogLevel = "warning" - LogLevelInfo LogLevel = "info" + // LogLevelInfo is info log level + LogLevelInfo LogLevel = "info" + // LogLevelVerbose is verbose log level LogLevelVerbose LogLevel = "verbose" - LogLevelDebug LogLevel = "debug" - LogLevelTrace LogLevel = "trace" + // LogLevelDebug is debug log level + LogLevelDebug LogLevel = "debug" + // LogLevelTrace is trace log level + LogLevelTrace LogLevel = "trace" ) // Options represents configuration for the ffmpeg command diff --git a/internal/pkg/mediaconvert/command_test.go b/internal/pkg/mediaconvert/command_test.go index 26beed3272..f8b36a7750 100644 --- a/internal/pkg/mediaconvert/command_test.go +++ b/internal/pkg/mediaconvert/command_test.go @@ -28,10 +28,11 @@ func Test_BuildVideoCommand_Scaling(t *testing.T) { Input: "pipe:0", UploadID: "1", Threads: 4, + LogLevel: mediaconvert.LogLevelDebug, ScalingLevels: []string{"720p", "480p"}, }) - const expected = `-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 -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8` + 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` require.Contains(t, expected, strings.Join(scaleCommand, " ")) } @@ -42,11 +43,12 @@ func Test_BuildVideoCommand_Scaling_NoRaw(t *testing.T) { Input: "pipe:0", UploadID: "1", Threads: 4, + LogLevel: mediaconvert.LogLevelDebug, Level: "720p", ScalingLevels: []string{"720p", "480p"}, }) - const expected = `-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 -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` require.Contains(t, expected, strings.Join(scaleCommand, " ")) } @@ -57,11 +59,12 @@ func Test_BuildVideoCommand_Raw_NoTranscode(t *testing.T) { Input: "pipe:0", UploadID: "1", Threads: 4, + LogLevel: mediaconvert.LogLevelDebug, Level: resconv.Level("651:490"), Transcode: false, }) - const expected = `"-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 -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, " ")) } @@ -72,11 +75,12 @@ func Test_BuildVideoCommand_Raw_Transcode(t *testing.T) { Input: "pipe:0", UploadID: "1", Threads: 4, + LogLevel: mediaconvert.LogLevelDebug, Level: resconv.Level("651:490"), Transcode: true, }) - const expected = `-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 -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` require.Contains(t, expected, strings.Join(rawCommand, " ")) } diff --git a/internal/pkg/mediaconvert/transcoder.go b/internal/pkg/mediaconvert/transcoder.go index 46b8206083..9667073bf0 100644 --- a/internal/pkg/mediaconvert/transcoder.go +++ b/internal/pkg/mediaconvert/transcoder.go @@ -44,6 +44,7 @@ type Transcoder struct { logger *zap.Logger } +// Command represents a ffmpeg command type Command struct { cmd *exec.Cmd stdoutBuf bytes.Buffer @@ -186,7 +187,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er if cmdErr != nil { logger.Error("can not create a new command", zap.Error(cmdErr), zap.Strings("args", args)) go uploader.Cancel() - return nil, errors.Wrapf(err, "can not create a new command") + return nil, errors.Wrapf(cmdErr, "can not create a new command") } var command = Command{ @@ -209,13 +210,19 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er logger.Debug("phase 7: wait for the result") for _, cmd := range cmds { - if err = cmd.cmd.Wait(); err != nil { - logger.Error("can not wait for command end ", zap.Error(err)) - os.Stdout.Write(cmd.stdoutBuf.Bytes()) - os.Stderr.Write(cmd.stderrBuf.Bytes()) - go uploader.Cancel() - return nil, errors.Wrapf(err, "can not wait for command end") + if err = cmd.cmd.Wait(); err == nil { + continue } + + logger.Error("can not wait for command end ", zap.Error(err)) + if _, err = os.Stdout.Write(cmd.stdoutBuf.Bytes()); err != nil { + logger.Error("can not write stdout ", zap.Error(err)) + } + if _, err = os.Stderr.Write(cmd.stderrBuf.Bytes()); err != nil { + logger.Error("can not write stderr", zap.Error(err)) + } + go uploader.Cancel() + return nil, errors.Wrapf(err, "can not wait for command end") } logger.Debug("phase 8: schedule cleanup") diff --git a/internal/pkg/storage/datalake.go b/internal/pkg/storage/datalake.go index 33d8b9e159..95c8755569 100644 --- a/internal/pkg/storage/datalake.go +++ b/internal/pkg/storage/datalake.go @@ -285,7 +285,7 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str defer func() { _ = file.Close() }() - if err := resp.BodyWriteTo(file); err != nil { + if err = resp.BodyWriteTo(file); err != nil { logger.Debug("can't write to file", zap.Error(err)) return err } diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index e14782f0a7..4e1515c80f 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -277,7 +277,7 @@ func (u *uploaderImpl) uploadAndDelete(f string) { return } - if err := waitFileExists(f); os.IsNotExist(err) { + if err := waitFileExists(f); err != nil { if os.IsNotExist(err) { logger.Debug("file does not exist", zap.Error(err)) } else { @@ -392,9 +392,10 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) { func waitFileExists(file string) error { var err error + var stat os.FileInfo for range 10 { - stat, err := os.Stat(file) + stat, err = os.Stat(file) if err == nil && stat.Size() > 0 { return nil }