From b35382b21679d563a808acd78c0c07fdc44c8c4f Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Fri, 30 May 2025 17:32:03 +0700 Subject: [PATCH 1/3] fix: properly initialize uploader Signed-off-by: Alexander Onnikov --- internal/pkg/uploader/uploader.go | 97 ++++++++++++++++++++----------- 1 file changed, 64 insertions(+), 33 deletions(-) diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index 77b7cfccda..247f5ec84a 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -91,20 +91,10 @@ func New(ctx context.Context, s storage.Storage, opts Options) Uploader { res.uploadCtx, res.uploadCancel = context.WithCancel(context.Background()) - _ = os.MkdirAll(opts.Dir, os.ModePerm) - res.workerWaitGroup.Add(1) - - go func() { - defer res.workerWaitGroup.Done() - initFiles, _ := os.ReadDir(opts.Dir) - for _, f := range initFiles { - var filePath = filepath.Join(opts.Dir, f.Name()) - if filePath == opts.SourceFile { - continue - } - res.filesCh <- filePath - } - }() + err := os.MkdirAll(opts.Dir, os.ModePerm) + if err != nil { + res.logger.Error("can not create upload directory", zap.Error(err), zap.String("dir", opts.Dir)) + } return res } @@ -117,6 +107,35 @@ func (u *uploaderImpl) Cancel() { u.stop(true) } +func (u *uploaderImpl) scanInitialFiles() { + u.workerWaitGroup.Add(1) + + go func() { + defer u.workerWaitGroup.Done() + + initFiles, err := os.ReadDir(u.options.Dir) + if err != nil { + u.logger.Error("failed to read initial files", zap.Error(err), zap.String("dir", u.options.Dir)) + return + } + + for _, f := range initFiles { + if f.IsDir() { + continue + } + + var filePath = filepath.Join(u.options.Dir, f.Name()) + if filePath == u.options.SourceFile { + continue + } + u.filesCh <- filePath + } + + u.logger.Info("initial file scan complete", zap.String("dir", u.options.Dir), zap.Int("count", len(initFiles))) + }() + +} + func (u *uploaderImpl) stop(rollback bool) { close(u.watcherStopCh) <-u.watcherDoneCh @@ -147,8 +166,14 @@ func (u *uploaderImpl) stop(rollback bool) { } func (u *uploaderImpl) Start() { + watcherReady := make(chan struct{}) + u.startWorkers() - u.startWatch() + go u.startWatch(watcherReady) + + <-watcherReady + + u.scanInitialFiles() } func (u *uploaderImpl) startWorkers() { @@ -230,13 +255,11 @@ func (u *uploaderImpl) uploadAndDelete(f string) { if err != nil { logger.Error("attempt failed", zap.Error(err)) } else { - if !u.shouldDeleteOnStop(f) { - _ = os.Remove(f) - logger.Debug("removed file locally") - } + // Mark the file as uploaded u.sentFiles.Store(f, struct{}{}) logger.Debug("file uploaded") + // Update the file's parent if SourceFile is set if u.options.SourceFile != "" { err = u.storage.SetParent(ctx, f, u.options.SourceFile) if err != nil { @@ -244,6 +267,16 @@ func (u *uploaderImpl) uploadAndDelete(f string) { } } + // Delete the file locally if it should be deleted + if !u.shouldDeleteOnStop(f) { + if err := os.Remove(f); err != nil { + logger.Error("failed to remove file locally", zap.Error(err), zap.String("file", f)) + } else { + logger.Debug("removed file locally") + } + } + + break } @@ -251,28 +284,26 @@ func (u *uploaderImpl) uploadAndDelete(f string) { } } -func (u *uploaderImpl) startWatch() { +func (u *uploaderImpl) startWatch(ready chan<- struct{}) { + defer close(u.watcherDoneCh) + var logger = u.logger.With(zap.String("func", "startWatch")) var watcher, err = inotify.NewWatcher() if err != nil { logger.Error("can not start file watcher", zap.Error(err)) + close(ready) + return + } + defer watcher.Close() + + if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite | inotifyMovedTo); err != nil { + logger.Error("can not start watching", zap.Error(err)) + close(ready) return } - if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite); err != nil { - logger.Error("can not start watching for close write", zap.Error(err)) - return - } - if err := watcher.AddWatch(u.options.Dir, inotifyMovedTo); err != nil { - logger.Error("can not start watching for moved to", zap.Error(err)) - return - } - defer func() { - _ = watcher.Close() - close(u.watcherDoneCh) - }() - + close(ready) logger.Debug("watching for file updates") defer logger.Debug("done") From d2f20a5acc7ec45ba1b27801e1de8f10c68876f2 Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Fri, 30 May 2025 18:31:07 +0700 Subject: [PATCH 2/3] fix: typo fixes Signed-off-by: Alexander Onnikov --- internal/pkg/api/v1/transcoding/handler.go | 6 +++++- internal/pkg/mediaconvert/command.go | 22 +++++++++++----------- internal/pkg/mediaconvert/command_test.go | 4 ++-- internal/pkg/mediaconvert/coordinator.go | 2 +- internal/pkg/mediaconvert/scheduler.go | 14 ++++++++++---- internal/pkg/mediaconvert/stream.go | 2 +- 6 files changed, 30 insertions(+), 20 deletions(-) diff --git a/internal/pkg/api/v1/transcoding/handler.go b/internal/pkg/api/v1/transcoding/handler.go index 67c4ca1fd0..c3a3ecfc92 100644 --- a/internal/pkg/api/v1/transcoding/handler.go +++ b/internal/pkg/api/v1/transcoding/handler.go @@ -60,7 +60,11 @@ func (t *trascodeHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { return } - t.scheduler.Schedule(&task) + if err := t.scheduler.Schedule(&task); err != nil { + w.WriteHeader(http.StatusTooManyRequests) + return + } + w.WriteHeader(http.StatusOK) } diff --git a/internal/pkg/mediaconvert/command.go b/internal/pkg/mediaconvert/command.go index ceb95c85de..29245eccb2 100644 --- a/internal/pkg/mediaconvert/command.go +++ b/internal/pkg/mediaconvert/command.go @@ -32,7 +32,7 @@ import ( // Options represents configuration for the ffmpeg command type Options struct { Input string - OuputDir string + OutputDir string ScalingLevels []string Level string Threads int @@ -56,7 +56,7 @@ func newFfmpegCommand(ctx context.Context, in io.Reader, args []string) (*exec.C return result, nil } -func buildCommonComamnd(opts *Options) []string { +func buildCommonCommand(opts *Options) []string { return []string{ "-threads", fmt.Sprint(opts.Threads), "-i", opts.Input, @@ -65,24 +65,24 @@ func buildCommonComamnd(opts *Options) []string { // BuildAudioCommand returns flags for getting the audio from the input func BuildAudioCommand(opts *Options) []string { - var commonPart = buildCommonComamnd(opts) + var commonPart = buildCommonCommand(opts) return append(commonPart, "-vn", "-acodec", - "copy", filepath.Join(opts.OuputDir, opts.UploadID), + "copy", filepath.Join(opts.OutputDir, opts.UploadID), ) } // BuildRawVideoCommand returns an extremely lightweight ffmpeg command for converting raw video without extra cost. func BuildRawVideoCommand(opts *Options) []string { - return append(buildCommonComamnd(opts), + return append(buildCommonCommand(opts), "-c:a", "copy", // Copy audio stream "-c:v", "copy", // Copy video stream "-hls_time", "5", "-hls_flags", "split_by_time", "-hls_list_size", "0", - "-hls_segment_filename", filepath.Join(opts.OuputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", opts.Level)), - filepath.Join(opts.OuputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, opts.Level))) + "-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))) } // BuildThumbnailCommand creates a command that creates a thumbnail for the input video @@ -90,13 +90,13 @@ func BuildThumbnailCommand(opts *Options) []string { return append([]string{}, "-i", opts.Input, "-vframes", "1", - filepath.Join(opts.OuputDir, opts.UploadID, opts.UploadID+".jpg"), + filepath.Join(opts.OutputDir, opts.UploadID, opts.UploadID+".jpg"), ) } // BuildScalingVideoCommand returns flags for ffmpeg for video scaling func BuildScalingVideoCommand(opts *Options) []string { - var result = buildCommonComamnd(opts) + var result = buildCommonCommand(opts) for _, level := range opts.ScalingLevels { result = append(result, @@ -109,8 +109,8 @@ func BuildScalingVideoCommand(opts *Options) []string { "-hls_time", "5", "-hls_flags", "split_by_time", "-hls_list_size", "0", - "-hls_segment_filename", filepath.Join(opts.OuputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", level)), - filepath.Join(opts.OuputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, level))) + "-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))) } return result diff --git a/internal/pkg/mediaconvert/command_test.go b/internal/pkg/mediaconvert/command_test.go index b67fe8060b..09ef87cd8a 100644 --- a/internal/pkg/mediaconvert/command_test.go +++ b/internal/pkg/mediaconvert/command_test.go @@ -24,7 +24,7 @@ import ( func Test_BuildVideoCommand_Scaling(t *testing.T) { var scaleCommand = mediaconvert.BuildScalingVideoCommand(&mediaconvert.Options{ - OuputDir: "test", + OutputDir: "test", Input: "pipe:0", UploadID: "1", Threads: 4, @@ -38,7 +38,7 @@ func Test_BuildVideoCommand_Scaling(t *testing.T) { func Test_BuildVideoCommand_Raw(t *testing.T) { var rawCommand = mediaconvert.BuildRawVideoCommand(&mediaconvert.Options{ - OuputDir: "test", + OutputDir: "test", Input: "pipe:0", UploadID: "1", Threads: 4, diff --git a/internal/pkg/mediaconvert/coordinator.go b/internal/pkg/mediaconvert/coordinator.go index 843716ce79..80c4b3b21a 100644 --- a/internal/pkg/mediaconvert/coordinator.go +++ b/internal/pkg/mediaconvert/coordinator.go @@ -97,7 +97,7 @@ func (s *StreamCoordinator) NewUpload(ctx context.Context, info handler.FileInfo var commandOptions = Options{ Input: "pipe:0", - OuputDir: s.conf.OutputDir, + OutputDir: s.conf.OutputDir, Threads: s.conf.MaxThreadCount, UploadID: info.ID, Level: level, diff --git a/internal/pkg/mediaconvert/scheduler.go b/internal/pkg/mediaconvert/scheduler.go index 5c686629ee..4c20d52ce6 100644 --- a/internal/pkg/mediaconvert/scheduler.go +++ b/internal/pkg/mediaconvert/scheduler.go @@ -62,15 +62,15 @@ type Scheduler struct { } // Schedule schedules a task to transcode -func (p *Scheduler) Schedule(t *Task) { +func (p *Scheduler) Schedule(t *Task) error { t.ID = uuid.NewString() t.Status = "planned" select { case p.taskCh <- t: - p.logger.Sugar().Debugf("task %v is scheduled", t) + return nil default: - p.logger.Error("task channel is full") + return fmt.Errorf("task queue is full") } } @@ -127,6 +127,12 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) { return } + defer func() { + if err := os.RemoveAll(destinationFolder); err != nil { + logger.Error("failed to cleanup temporary folder", zap.Error(err)) + } + }() + logger.Debug("phase 3: get the remote file") remoteStorage, err := storage.NewStorageByURL(ctx, p.cfg.Endpoint(), p.cfg.EndpointURL.Scheme, tokenString, task.Workspace) @@ -176,7 +182,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) { var level = resconv.Level(res) var opts = Options{ Input: sourceFilePath, - OuputDir: p.cfg.OutputDir, + OutputDir: p.cfg.OutputDir, Level: level, ScalingLevels: append(resconv.SubLevels(res), level), UploadID: task.ID, diff --git a/internal/pkg/mediaconvert/stream.go b/internal/pkg/mediaconvert/stream.go index ab47c8e771..db207e9f61 100644 --- a/internal/pkg/mediaconvert/stream.go +++ b/internal/pkg/mediaconvert/stream.go @@ -115,7 +115,7 @@ func (s *StreamCoordinator) AsConcatableUpload(upload handler.Upload) handler.Co func (w *Stream) start(ctx context.Context, options *Options) error { defer w.logger.Debug("start done") w.reader = w.writer.Transpile() - if err := manifest.GenerateHLSPlaylist(append(options.ScalingLevels, options.Level), options.OuputDir, options.UploadID); err != nil { + if err := manifest.GenerateHLSPlaylist(append(options.ScalingLevels, options.Level), options.OutputDir, options.UploadID); err != nil { return err } w.commandGroup.Add(1) From faeca99748169d1f7663ca20901c1de5c42775a2 Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Fri, 30 May 2025 20:39:28 +0700 Subject: [PATCH 3/3] ci fixes Signed-off-by: Alexander Onnikov --- .golangci.yaml | 2 +- internal/pkg/mediaconvert/command_test.go | 10 +++++----- internal/pkg/mediaconvert/coordinator.go | 2 +- internal/pkg/mediaconvert/scheduler.go | 4 ++-- internal/pkg/uploader/uploader.go | 10 ++++++---- 5 files changed, 15 insertions(+), 13 deletions(-) diff --git a/.golangci.yaml b/.golangci.yaml index 16f9ae9d7b..5fd9dcf2b3 100644 --- a/.golangci.yaml +++ b/.golangci.yaml @@ -61,7 +61,7 @@ linters-settings: dupl: threshold: 150 funlen: - lines: 160 + lines: 180 statements: 100 goconst: min-len: 2 diff --git a/internal/pkg/mediaconvert/command_test.go b/internal/pkg/mediaconvert/command_test.go index 09ef87cd8a..b9cc151a77 100644 --- a/internal/pkg/mediaconvert/command_test.go +++ b/internal/pkg/mediaconvert/command_test.go @@ -24,7 +24,7 @@ import ( func Test_BuildVideoCommand_Scaling(t *testing.T) { var scaleCommand = mediaconvert.BuildScalingVideoCommand(&mediaconvert.Options{ - OutputDir: "test", + OutputDir: "test", Input: "pipe:0", UploadID: "1", Threads: 4, @@ -39,10 +39,10 @@ func Test_BuildVideoCommand_Scaling(t *testing.T) { func Test_BuildVideoCommand_Raw(t *testing.T) { var rawCommand = mediaconvert.BuildRawVideoCommand(&mediaconvert.Options{ OutputDir: "test", - Input: "pipe:0", - UploadID: "1", - Threads: 4, - Level: resconv.Level("651:490"), + Input: "pipe:0", + UploadID: "1", + Threads: 4, + Level: resconv.Level("651:490"), }) const expected = `"-threads 4 -i pipe:0 -c:a copy -c:v copy -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` diff --git a/internal/pkg/mediaconvert/coordinator.go b/internal/pkg/mediaconvert/coordinator.go index 80c4b3b21a..ede68b3b65 100644 --- a/internal/pkg/mediaconvert/coordinator.go +++ b/internal/pkg/mediaconvert/coordinator.go @@ -97,7 +97,7 @@ func (s *StreamCoordinator) NewUpload(ctx context.Context, info handler.FileInfo var commandOptions = Options{ Input: "pipe:0", - OutputDir: s.conf.OutputDir, + OutputDir: s.conf.OutputDir, Threads: s.conf.MaxThreadCount, UploadID: info.ID, Level: level, diff --git a/internal/pkg/mediaconvert/scheduler.go b/internal/pkg/mediaconvert/scheduler.go index 4c20d52ce6..31507f1747 100644 --- a/internal/pkg/mediaconvert/scheduler.go +++ b/internal/pkg/mediaconvert/scheduler.go @@ -128,7 +128,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) { } defer func() { - if err := os.RemoveAll(destinationFolder); err != nil { + if err = os.RemoveAll(destinationFolder); err != nil { logger.Error("failed to cleanup temporary folder", zap.Error(err)) } }() @@ -182,7 +182,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) { var level = resconv.Level(res) var opts = Options{ Input: sourceFilePath, - OutputDir: p.cfg.OutputDir, + OutputDir: p.cfg.OutputDir, Level: level, ScalingLevels: append(resconv.SubLevels(res), level), UploadID: task.ID, diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index 247f5ec84a..ad1798fa00 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -133,7 +133,6 @@ func (u *uploaderImpl) scanInitialFiles() { u.logger.Info("initial file scan complete", zap.String("dir", u.options.Dir), zap.Int("count", len(initFiles))) }() - } func (u *uploaderImpl) stop(rollback bool) { @@ -276,7 +275,6 @@ func (u *uploaderImpl) uploadAndDelete(f string) { } } - break } @@ -295,9 +293,13 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) { close(ready) return } - defer watcher.Close() + defer func() { + if err := watcher.Close(); err != nil { + logger.Error("can not close watcher", zap.Error(err)) + } + }() - if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite | inotifyMovedTo); err != nil { + if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite|inotifyMovedTo); err != nil { logger.Error("can not start watching", zap.Error(err)) close(ready) return