From 36fe1deb355ffef326af2543676ee7ea5efc35b4 Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Thu, 23 Oct 2025 12:10:38 +0700 Subject: [PATCH] Do not transcode while recording Signed-off-by: Alexander Onnikov --- internal/pkg/mediaconvert/coordinator.go | 29 +------ internal/pkg/mediaconvert/stream.go | 101 ++++------------------- 2 files changed, 18 insertions(+), 112 deletions(-) diff --git a/internal/pkg/mediaconvert/coordinator.go b/internal/pkg/mediaconvert/coordinator.go index cfad9916d4..bd93d8a61f 100644 --- a/internal/pkg/mediaconvert/coordinator.go +++ b/internal/pkg/mediaconvert/coordinator.go @@ -23,7 +23,6 @@ import ( "strconv" "strings" "sync" - "sync/atomic" "time" "github.com/pkg/errors" @@ -47,8 +46,6 @@ type StreamCoordinator struct { outputDir string uploadOptions uploader.Options - activeTranscoding int32 - mainContext context.Context logger *zap.Logger @@ -96,12 +93,6 @@ func (s *StreamCoordinator) NewUpload(ctx context.Context, info handler.FileInfo done: make(chan struct{}), } - if atomic.AddInt32(&s.activeTranscoding, 1) > int32(s.conf.MaxParallelTranscodingCount) { - s.logger.Debug("run out of resources for scaling") - // atomic.AddInt32(&s.activeTranscoding, -1) - // TODO do not transcode - } - width, err := strconv.Atoi(info.MetaData["width"]) if err != nil { return nil, errors.Wrapf(err, "can not parse video width: %v", info.MetaData["width"]) @@ -118,15 +109,6 @@ func (s *StreamCoordinator) NewUpload(ctx context.Context, info handler.FileInfo Codec: extractCodec(info.MetaData["contentType"]), ContentType: extractContentType(info.MetaData["contentType"]), } - profiles := FastTranscodingProfiles(meta) - - var commandOptions = Options{ - Input: "pipe:0", - OutputDir: s.outputDir, - Threads: s.conf.MaxThreadCount, - UploadID: info.ID, - Profiles: profiles, - } if s.conf.EndpointURL != nil { s.logger.Sugar().Debugf("initializing uploader for %v", info) @@ -153,22 +135,13 @@ func (s *StreamCoordinator) NewUpload(ctx context.Context, info handler.FileInfo } stream.multipart = multipart } - // uploader for processed outputs - var contentUploader = uploader.New(s.mainContext, stg, opts) - stream.contentUploader = contentUploader } s.streams.Store(stream.info.ID, stream) - if err := stream.start(s.mainContext, &commandOptions); err != nil { + if err := stream.start(s.mainContext); err != nil { return nil, err } - go func() { - stream.commandGroup.Wait() - atomic.AddInt32(&s.activeTranscoding, -1) - close(stream.done) - }() - s.manageTimeout(stream) s.logger.Debug("NewUpload", zap.String("done", info.ID)) diff --git a/internal/pkg/mediaconvert/stream.go b/internal/pkg/mediaconvert/stream.go index 95524361c9..fbb671c911 100644 --- a/internal/pkg/mediaconvert/stream.go +++ b/internal/pkg/mediaconvert/stream.go @@ -17,31 +17,26 @@ package mediaconvert import ( "context" "io" - "os/exec" "sync" "github.com/pkg/errors" - "github.com/hcengineering/stream/internal/pkg/manifest" "github.com/hcengineering/stream/internal/pkg/sharedpipe" "github.com/hcengineering/stream/internal/pkg/storage" - "github.com/hcengineering/stream/internal/pkg/uploader" "github.com/tus/tusd/v2/pkg/handler" "go.uber.org/zap" ) // Stream manages client's input and transcodes it based on the passed configuration type Stream struct { - contentUploader uploader.Uploader - logger *zap.Logger - info handler.FileInfo - writer *sharedpipe.Writer - reader *sharedpipe.Reader - storage storage.Storage - multipart *MultipartUpload + logger *zap.Logger + info handler.FileInfo + writer *sharedpipe.Writer + reader *sharedpipe.Reader + storage storage.Storage + multipart *MultipartUpload - commandGroup sync.WaitGroup - done chan struct{} + done chan struct{} } var _ handler.Upload = (*Stream)(nil) @@ -108,16 +103,6 @@ func (w *Stream) Terminate(ctx context.Context) error { var wg sync.WaitGroup - // cancel upload if in progress - if w.contentUploader != nil { - wg.Add(1) - go func() { - defer wg.Done() - w.commandGroup.Wait() - w.contentUploader.Cancel() - }() - } - // cancel multipart upload if in progress if w.multipart != nil { wg.Add(1) @@ -131,6 +116,9 @@ func (w *Stream) Terminate(ctx context.Context) error { wg.Wait() + // Signal that the stream is done + close(w.done) + return nil } @@ -154,15 +142,7 @@ func (w *Stream) FinishUpload(ctx context.Context) error { } var wg sync.WaitGroup - - if w.contentUploader != nil { - wg.Add(1) - go func() { - defer wg.Done() - w.commandGroup.Wait() - w.contentUploader.Stop() - }() - } + var completeErr error // finalize raw multipart stream if supported if w.multipart != nil { @@ -171,30 +151,18 @@ func (w *Stream) FinishUpload(ctx context.Context) error { defer wg.Done() if err := w.multipart.Complete(ctx); err != nil { w.logger.Error("multipart upload complete failed", zap.Error(err)) + completeErr = err return } - - if metaProvider, ok := w.storage.(storage.MetaProvider); ok { - metaErr := metaProvider.PatchMeta( - ctx, - w.info.ID, - &storage.Metadata{ - "hls": map[string]any{ - "source": manifest.MasterPlaylistFileName(w.info.ID), - "thumbnail": manifest.ThumbnailFileName(w.info.ID), - }, - }, - ) - if metaErr != nil { - w.logger.Error("can not patch the source file", zap.Error(metaErr)) - } - } }() } wg.Wait() - return nil + // Signal that the stream is done + close(w.done) + + return completeErr } // AsConcatableUpload returns tusd handler.ConcatableUpload @@ -203,43 +171,8 @@ func (s *StreamCoordinator) AsConcatableUpload(upload handler.Upload) handler.Co return upload.(*Stream) } -func (w *Stream) start(ctx context.Context, options *Options) error { +func (w *Stream) start(ctx context.Context) error { defer w.logger.Debug("start done") w.reader = w.writer.Transpile() - if err := manifest.GenerateHLSPlaylist(options.Profiles, options.OutputDir, options.UploadID); err != nil { - return err - } - - var argsSlice = [][]string{ - BuildThumbnailCommand(options), - BuildVideoCommand(options), - } - - var cmds []*exec.Cmd - for idx, args := range argsSlice { - reader := w.reader - if idx > 0 { - reader = w.writer.Transpile() - } - - cmd, cmdErr := newFfmpegCommand(ctx, reader, args) - if cmdErr != nil { - w.logger.Error("can not create a new command", zap.Error(cmdErr), zap.Strings("args", args)) - return errors.Wrapf(cmdErr, "can not create a new command") - } - cmds = append(cmds, cmd) - } - - w.commandGroup.Add(1) - go func() { - defer w.commandGroup.Done() - executor := NewCommandExecutor(ctx) - if execErr := executor.Execute(cmds); execErr != nil { - w.logger.Error("can not execute command", zap.Error(execErr)) - } - }() - - go w.contentUploader.Start() - return nil }