diff --git a/.github/workflows/docker-push.yaml b/.github/workflows/docker-push.yaml index 0308ce654b..ad67bcee7c 100644 --- a/.github/workflows/docker-push.yaml +++ b/.github/workflows/docker-push.yaml @@ -41,4 +41,4 @@ jobs: context: . platforms: linux/amd64,linux/arm64 push: true - tags: ${{ inputs.version }} + tags: hardcoreeng/huly-stream:${{ inputs.version }} diff --git a/cmd/huly-stream/main.go b/cmd/huly-stream/main.go index 2730cd8b35..f660a2511f 100644 --- a/cmd/huly-stream/main.go +++ b/cmd/huly-stream/main.go @@ -32,7 +32,7 @@ import ( tusd "github.com/tus/tusd/v2/pkg/handler" ) -const basePath = "/transcoding" +const basePath = "/recording" func main() { var ctx, cancel = signal.NotifyContext( @@ -43,16 +43,13 @@ func main() { syscall.SIGQUIT, ) defer cancel() - ctx = log.WithLoggerFields(ctx) var logger = log.FromContext(ctx) var conf = must(config.FromEnv()) - logger.Sugar().Debugf("provided config is %v", conf) mustNoError(os.MkdirAll(conf.OutputDir, os.ModePerm)) - if conf.PprofEnabled { go pprof.ListenAndServe(ctx, "localhost:6060") } @@ -71,10 +68,12 @@ func main() { Logger: slog.New(slog.NewTextHandler(discardTextHandler{}, nil)), })) - http.Handle("/transcoding/", http.StripPrefix("/transcoding/", handler)) - http.Handle("/transcoding", http.StripPrefix("/transcoding", handler)) + http.Handle("/recording/", http.StripPrefix("/recording/", handler)) + http.Handle("/recording", http.StripPrefix("/recording", handler)) go func() { + logger.Info("started to listen") + defer logger.Info("server has finished") // #nosec var err = http.ListenAndServe(conf.ServeURL, nil) if err != nil { diff --git a/go.mod b/go.mod index 89419cb728..3775572bbf 100644 --- a/go.mod +++ b/go.mod @@ -3,42 +3,42 @@ module github.com/huly-stream go 1.23.2 require ( - github.com/aws/aws-sdk-go-v2 v1.32.3 - github.com/aws/aws-sdk-go-v2/config v1.28.1 - github.com/aws/aws-sdk-go-v2/credentials v1.17.42 - github.com/aws/aws-sdk-go-v2/service/s3 v1.66.2 - github.com/aws/smithy-go v1.22.0 + github.com/aws/aws-sdk-go-v2 v1.36.1 + github.com/aws/aws-sdk-go-v2/config v1.29.6 + 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/aws/smithy-go v1.22.3 github.com/fsnotify/fsnotify v1.8.0 github.com/google/uuid v1.6.0 github.com/kelseyhightower/envconfig v1.4.0 github.com/pkg/errors v0.9.1 - github.com/stretchr/testify v1.9.0 + github.com/stretchr/testify v1.10.0 github.com/tus/tusd/v2 v2.6.0 - github.com/valyala/fasthttp v1.58.0 + github.com/valyala/fasthttp v1.59.0 go.uber.org/zap v1.27.0 - golang.org/x/exp v0.0.0-20230626212559-97b1e661b5df + golang.org/x/exp v0.0.0-20250215185904-eff6e970281f ) require ( github.com/andybalholm/brotli v1.1.1 // indirect - github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.6 // indirect - github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.18 // indirect - github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.22 // indirect - github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.22 // indirect - github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1 // indirect - github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.22 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.0 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.3 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.3 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.3 // indirect - github.com/aws/aws-sdk-go-v2/service/sso v1.24.3 // indirect - github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.3 // indirect - github.com/aws/aws-sdk-go-v2/service/sts v1.32.3 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.9 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.28 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.32 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.32 // indirect + github.com/aws/aws-sdk-go-v2/internal/ini v1.8.2 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.32 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.2 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.6.0 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.13 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.13 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.24.15 // indirect + github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.14 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.33.14 // indirect github.com/davecgh/go-spew v1.1.1 // indirect github.com/klauspost/compress v1.17.11 // indirect github.com/pmezard/go-difflib v1.0.0 // indirect github.com/valyala/bytebufferpool v1.0.0 // indirect go.uber.org/multierr v1.11.0 // indirect - golang.org/x/sys v0.27.0 // indirect + golang.org/x/sys v0.30.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index 11e9c85d56..1ffa95ca7c 100644 --- a/go.sum +++ b/go.sum @@ -2,42 +2,42 @@ github.com/Acconut/go-httptest-recorder v1.0.0 h1:TAv2dfnqp/l+SUvIaMAUK4GeN4+wqb github.com/Acconut/go-httptest-recorder v1.0.0/go.mod h1:CwQyhTH1kq/gLyWiRieo7c0uokpu3PXeyF/nZjUNtmM= github.com/andybalholm/brotli v1.1.1 h1:PR2pgnyFznKEugtsUo0xLdDop5SKXd5Qf5ysW+7XdTA= github.com/andybalholm/brotli v1.1.1/go.mod h1:05ib4cKhjx3OQYUY22hTVd34Bc8upXjOLL2rKwwZBoA= -github.com/aws/aws-sdk-go-v2 v1.32.3 h1:T0dRlFBKcdaUPGNtkBSwHZxrtis8CQU17UpNBZYd0wk= -github.com/aws/aws-sdk-go-v2 v1.32.3/go.mod h1:2SK5n0a2karNTv5tbP1SjsX0uhttou00v/HpXKM1ZUo= -github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.6 h1:pT3hpW0cOHRJx8Y0DfJUEQuqPild8jRGmSFmBgvydr0= -github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.6/go.mod h1:j/I2++U0xX+cr44QjHay4Cvxj6FUbnxrgmqN3H1jTZA= -github.com/aws/aws-sdk-go-v2/config v1.28.1 h1:oxIvOUXy8x0U3fR//0eq+RdCKimWI900+SV+10xsCBw= -github.com/aws/aws-sdk-go-v2/config v1.28.1/go.mod h1:bRQcttQJiARbd5JZxw6wG0yIK3eLeSCPdg6uqmmlIiI= -github.com/aws/aws-sdk-go-v2/credentials v1.17.42 h1:sBP0RPjBU4neGpIYyx8mkU2QqLPl5u9cmdTWVzIpHkM= -github.com/aws/aws-sdk-go-v2/credentials v1.17.42/go.mod h1:FwZBfU530dJ26rv9saAbxa9Ej3eF/AK0OAY86k13n4M= -github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.18 h1:68jFVtt3NulEzojFesM/WVarlFpCaXLKaBxDpzkQ9OQ= -github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.18/go.mod h1:Fjnn5jQVIo6VyedMc0/EhPpfNlPl7dHV916O6B+49aE= -github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.22 h1:Jw50LwEkVjuVzE1NzkhNKkBf9cRN7MtE1F/b2cOKTUM= -github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.22/go.mod h1:Y/SmAyPcOTmpeVaWSzSKiILfXTVJwrGmYZhcRbhWuEY= -github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.22 h1:981MHwBaRZM7+9QSR6XamDzF/o7ouUGxFzr+nVSIhrs= -github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.22/go.mod h1:1RA1+aBEfn+CAB/Mh0MB6LsdCYCnjZm7tKXtnk499ZQ= -github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1 h1:VaRN3TlFdd6KxX1x3ILT5ynH6HvKgqdiXoTxAF4HQcQ= -github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1/go.mod h1:FbtygfRFze9usAadmnGJNc8KsP346kEe+y2/oyhGAGc= -github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.22 h1:yV+hCAHZZYJQcwAaszoBNwLbPItHvApxT0kVIw6jRgs= -github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.22/go.mod h1:kbR1TL8llqB1eGnVbybcA4/wgScxdylOdyAd51yxPdw= -github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.0 h1:TToQNkvGguu209puTojY/ozlqy2d/SFNcoLIqTFi42g= -github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.0/go.mod h1:0jp+ltwkf+SwG2fm/PKo8t4y8pJSgOCO4D8Lz3k0aHQ= -github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.3 h1:kT6BcZsmMtNkP/iYMcRG+mIEA/IbeiUimXtGmqF39y0= -github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.3/go.mod h1:Z8uGua2k4PPaGOYn66pK02rhMrot3Xk3tpBuUFPomZU= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.3 h1:qcxX0JYlgWH3hpPUnd6U0ikcl6LLA9sLkXE2w1fpMvY= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.3/go.mod h1:cLSNEmI45soc+Ef8K/L+8sEA3A3pYFEYf5B5UI+6bH4= -github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.3 h1:ZC7Y/XgKUxwqcdhO5LE8P6oGP1eh6xlQReWNKfhvJno= -github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.3/go.mod h1:WqfO7M9l9yUAw0HcHaikwRd/H6gzYdz7vjejCA5e2oY= -github.com/aws/aws-sdk-go-v2/service/s3 v1.66.2 h1:p9TNFL8bFUMd+38YIpTAXpoxyz0MxC7FlbFEH4P4E1U= -github.com/aws/aws-sdk-go-v2/service/s3 v1.66.2/go.mod h1:fNjyo0Coen9QTwQLWeV6WO2Nytwiu+cCcWaTdKCAqqE= -github.com/aws/aws-sdk-go-v2/service/sso v1.24.3 h1:UTpsIf0loCIWEbrqdLb+0RxnTXfWh2vhw4nQmFi4nPc= -github.com/aws/aws-sdk-go-v2/service/sso v1.24.3/go.mod h1:FZ9j3PFHHAR+w0BSEjK955w5YD2UwB/l/H0yAK3MJvI= -github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.3 h1:2YCmIXv3tmiItw0LlYf6v7gEHebLY45kBEnPezbUKyU= -github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.3/go.mod h1:u19stRyNPxGhj6dRm+Cdgu6N75qnbW7+QN0q0dsAk58= -github.com/aws/aws-sdk-go-v2/service/sts v1.32.3 h1:wVnQ6tigGsRqSWDEEyH6lSAJ9OyFUsSnbaUWChuSGzs= -github.com/aws/aws-sdk-go-v2/service/sts v1.32.3/go.mod h1:VZa9yTFyj4o10YGsmDO4gbQJUvvhY72fhumT8W4LqsE= -github.com/aws/smithy-go v1.22.0 h1:uunKnWlcoL3zO7q+gG2Pk53joueEOsnNB28QdMsmiMM= -github.com/aws/smithy-go v1.22.0/go.mod h1:irrKGvNn1InZwb2d7fkIRNucdfwR8R+Ts3wxYa/cJHg= +github.com/aws/aws-sdk-go-v2 v1.36.1 h1:iTDl5U6oAhkNPba0e1t1hrwAo02ZMqbrGq4k5JBWM5E= +github.com/aws/aws-sdk-go-v2 v1.36.1/go.mod h1:5PMILGVKiW32oDzjj6RU52yrNrDPUHcbZQYr1sM7qmM= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.9 h1:VZPDrbzdsU1ZxhyWrvROqLY0nxFWgMCAzhn/nYz3X48= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.9/go.mod h1:3XkePX5dSaxveLAYY7nsbsZZrKxCyEuE5pM4ziFxyGg= +github.com/aws/aws-sdk-go-v2/config v1.29.6 h1:fqgqEKK5HaZVWLQoLiC9Q+xDlSp+1LYidp6ybGE2OGg= +github.com/aws/aws-sdk-go-v2/config v1.29.6/go.mod h1:Ft+WLODzDQmCTHDvqAH1JfC2xxbZ0MxpZAcJqmE1LTQ= +github.com/aws/aws-sdk-go-v2/credentials v1.17.59 h1:9btwmrt//Q6JcSdgJOLI98sdr5p7tssS9yAsGe8aKP4= +github.com/aws/aws-sdk-go-v2/credentials v1.17.59/go.mod h1:NM8fM6ovI3zak23UISdWidyZuI1ghNe2xjzUZAyT+08= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.28 h1:KwsodFKVQTlI5EyhRSugALzsV6mG/SGrdjlMXSZSdso= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.28/go.mod h1:EY3APf9MzygVhKuPXAc5H+MkGb8k/DOSQjWS0LgkKqI= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.32 h1:BjUcr3X3K0wZPGFg2bxOWW3VPN8rkE3/61zhP+IHviA= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.32/go.mod h1:80+OGC/bgzzFFTUmcuwD0lb4YutwQeKLFpmt6hoWapU= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.32 h1:m1GeXHVMJsRsUAqG6HjZWx9dj7F5TR+cF1bjyfYyBd4= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.32/go.mod h1:IitoQxGfaKdVLNg0hD8/DXmAqNy0H4K2H2Sf91ti8sI= +github.com/aws/aws-sdk-go-v2/internal/ini v1.8.2 h1:Pg9URiobXy85kgFev3og2CuOZ8JZUBENF+dcgWBaYNk= +github.com/aws/aws-sdk-go-v2/internal/ini v1.8.2/go.mod h1:FbtygfRFze9usAadmnGJNc8KsP346kEe+y2/oyhGAGc= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.32 h1:OIHj/nAhVzIXGzbAE+4XmZ8FPvro3THr6NlqErJc3wY= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.32/go.mod h1:LiBEsDo34OJXqdDlRGsilhlIiXR7DL+6Cx2f4p1EgzI= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.2 h1:D4oz8/CzT9bAEYtVhSBmFj2dNOtaHOtMKc2vHBwYizA= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.2/go.mod h1:Za3IHqTQ+yNcRHxu1OFucBh0ACZT4j4VQFF0BqpZcLY= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.6.0 h1:kT2WeWcFySdYpPgyqJMSUE7781Qucjtn6wBvrgm9P+M= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.6.0/go.mod h1:WYH1ABybY7JK9TITPnk6ZlP7gQB8psI4c9qDmMsnLSA= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.13 h1:SYVGSFQHlchIcy6e7x12bsrxClCXSP5et8cqVhL8cuw= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.13/go.mod h1:kizuDaLX37bG5WZaoxGPQR/LNFXpxp0vsUnqfkWXfNE= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.13 h1:OBsrtam3rk8NfBEq7OLOMm5HtQ9Yyw32X4UQMya/wjw= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.13/go.mod h1:3U4gFA5pmoCOja7aq4nSaIAGbaOHv2Yl2ug018cmC+Q= +github.com/aws/aws-sdk-go-v2/service/s3 v1.77.0 h1:RCOi1rDmLqOICym/6UeS2cqKED4T4m966w2rl1HfL+g= +github.com/aws/aws-sdk-go-v2/service/s3 v1.77.0/go.mod h1:VC4EKSHqT3nzOcU955VWHMGsQ+w67wfAUBSjC8NOo8U= +github.com/aws/aws-sdk-go-v2/service/sso v1.24.15 h1:/eE3DogBjYlvlbhd2ssWyeuovWunHLxfgw3s/OJa4GQ= +github.com/aws/aws-sdk-go-v2/service/sso v1.24.15/go.mod h1:2PCJYpi7EKeA5SkStAmZlF6fi0uUABuhtF8ILHjGc3Y= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.14 h1:M/zwXiL2iXUrHputuXgmO94TVNmcenPHxgLXLutodKE= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.14/go.mod h1:RVwIw3y/IqxC2YEXSIkAzRDdEU1iRabDPaYjpGCbCGQ= +github.com/aws/aws-sdk-go-v2/service/sts v1.33.14 h1:TzeR06UCMUq+KA3bDkujxK1GVGy+G8qQN/QVYzGLkQE= +github.com/aws/aws-sdk-go-v2/service/sts v1.33.14/go.mod h1:dspXf/oYWGWo6DEvj98wpaTeqt5+DMidZD0A9BYTizc= +github.com/aws/smithy-go v1.22.3 h1:Z//5NuZCSW6R4PhQ93hShNbyBbn8BWCmCVCt+Q8Io5k= +github.com/aws/smithy-go v1.22.3/go.mod h1:t1ufH5HMublsJYulve2RKmHDC15xu1f26kHCp/HgceI= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/fsnotify/fsnotify v1.8.0 h1:dAwr6QBTBZIkG8roQaJjGof0pp0EeF+tNV7YBP3F/8M= @@ -54,14 +54,14 @@ github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg= -github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/tus/tusd/v2 v2.6.0 h1:Je243QDKnFTvm/WkLH2bd1oQ+7trolrflRWyuI0PdWI= github.com/tus/tusd/v2 v2.6.0/go.mod h1:1Eb1lBoSRBfYJ/mQfFVjyw8ZdNMdBqW17vgQKl3Ah9g= github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= -github.com/valyala/fasthttp v1.58.0 h1:GGB2dWxSbEprU9j0iMJHgdKYJVDyjrOwF9RE59PbRuE= -github.com/valyala/fasthttp v1.58.0/go.mod h1:SYXvHHaFp7QZHGKSHmoMipInhrI5StHrhDTYVEjK/Kw= +github.com/valyala/fasthttp v1.59.0 h1:Qu0qYHfXvPk1mSLNqcFtEk6DpxgA26hy6bmydotDpRI= +github.com/valyala/fasthttp v1.59.0/go.mod h1:GTxNb9Bc6r2a9D0TWNSPwDz78UxnTGBViY3xZNEqyYU= github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU= github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= @@ -70,14 +70,14 @@ go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= -golang.org/x/exp v0.0.0-20230626212559-97b1e661b5df h1:UA2aFVmmsIlefxMk29Dp2juaUSth8Pyn3Tq5Y5mJGME= -golang.org/x/exp v0.0.0-20230626212559-97b1e661b5df/go.mod h1:FXUEEKJgO7OQYeo8N01OfiKP8RXMtf6e8aTskBGqWdc= -golang.org/x/net v0.31.0 h1:68CPQngjLL0r2AlUKiSxtQFKvzRVbnzLwMUn5SzcLHo= -golang.org/x/net v0.31.0/go.mod h1:P4fl1q7dY2hnZFxEk4pPSkDHF+QqjitcnDjUQyMM+pM= -golang.org/x/sys v0.27.0 h1:wBqf8DvsY9Y/2P8gAfPDEYNuS30J4lPHJxXSb/nJZ+s= -golang.org/x/sys v0.27.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= -golang.org/x/text v0.20.0 h1:gK/Kv2otX8gz+wn7Rmb3vT96ZwuoxnQlY+HlJVj7Qug= -golang.org/x/text v0.20.0/go.mod h1:D4IsuqiFMhST5bX19pQ9ikHC2GsaKyk/oF+pn3ducp4= +golang.org/x/exp v0.0.0-20250215185904-eff6e970281f h1:oFMYAjX0867ZD2jcNiLBrI9BdpmEkvPyi5YrBGXbamg= +golang.org/x/exp v0.0.0-20250215185904-eff6e970281f/go.mod h1:BHOTPb3L19zxehTsLoJXVaTktb06DFgmdW6Wb9s8jqk= +golang.org/x/net v0.35.0 h1:T5GQRQb2y08kTAByq9L4/bz8cipCdA8FbRTXewonqY8= +golang.org/x/net v0.35.0/go.mod h1:EglIi67kWsHKlRzzVMUD93VMSWGFOMSZgxFjparz1Qk= +golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc= +golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/text v0.22.0 h1:bofq7m3/HAFvbF51jz3Q9wLg3jkvSPuiZu/pD1XwgtM= +golang.org/x/text v0.22.0/go.mod h1:YRoo4H8PVmsu+E3Ou7cqLVH8oXWIHVoX0jqUWALQhfY= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= diff --git a/internal/pkg/config/config.go b/internal/pkg/config/config.go index 1ef90f2192..9c296309c6 100644 --- a/internal/pkg/config/config.go +++ b/internal/pkg/config/config.go @@ -16,23 +16,24 @@ package config import ( "net/url" + "time" "github.com/kelseyhightower/envconfig" ) // Config represents configuration for the huly-stream application. type Config struct { - SecretToken string `split_words:"true" desc:"secret token for authorize requests"` - LogLevel string `split_words:"true" default:"debug" desc:"sets log level for the application"` - PprofEnabled bool `default:"false" split_words:"true" desc:"starts profile server on localhost:6060 if true"` - Insecure bool `default:"false" desc:"ignores authorization check if true"` - ServeURL string `split_words:"true" desc:"app listen url" default:"0.0.0.0:1080"` - EndpointURL *url.URL `split_words:"true" desc:"S3 or Datalake endpoint, example: s3://my-ip-address, datalake://my-ip-address"` - MaxCapacity int64 `split_words:"true" default:"6220800" desc:"represents the amount of maximum possible capacity for the transcoding. The default value is 1920 * 1080 * 3."` - MaxThreads int `split_words:"true" default:"4" desc:"means upper bound for the transcoing provider."` - OutputDir string `split_words:"true" default:"/tmp/transcoing/" desc:"path to the directory with transcoding result."` - RemoveContentOnUpload bool `split_words:"true" default:"true" desc:"deletes all content when content delivered if true"` - UploadRawContent bool `split_words:"true" default:"false" desc:"uploads content in raw quality to the endpoint if true"` + SecretToken string `split_words:"true" desc:"secret token for authorize requests"` + LogLevel string `split_words:"true" default:"debug" desc:"sets log level for the application"` + PprofEnabled bool `default:"false" split_words:"true" desc:"starts profile server on localhost:6060 if true"` + Insecure bool `default:"false" desc:"ignores authorization check if true"` + ServeURL string `split_words:"true" desc:"app listen url" default:"0.0.0.0:1080"` + EndpointURL *url.URL `split_words:"true" desc:"S3 or Datalake endpoint, example: s3://my-ip-address, datalake://my-ip-address"` + AuthURL *url.URL `split_words:"true" desc:"url to auth the upload"` + MaxCapacity int64 `split_words:"true" default:"6220800" desc:"represents the amount of maximum possible capacity for the transcoding. The default value is 1920 * 1080 * 3."` + MaxThreads int `split_words:"true" default:"4" desc:"means upper bound for the transcoing provider."` + OutputDir string `split_words:"true" default:"/tmp/transcoing/" desc:"path to the directory with transcoding result."` + Timeout time.Duration `default:"5m" desc:"timeout for the upload"` } // FromEnv creates new Config from env @@ -47,5 +48,9 @@ func FromEnv() (*Config, error) { return nil, err } + if *result.EndpointURL == (url.URL{}) { + result.EndpointURL = nil + } + return &result, nil } diff --git a/internal/pkg/manifest/hls.go b/internal/pkg/manifest/hls.go index 712e97f072..d8dd7a5361 100644 --- a/internal/pkg/manifest/hls.go +++ b/internal/pkg/manifest/hls.go @@ -15,111 +15,45 @@ package manifest import ( - "bufio" "fmt" - "strconv" + "os" + "path/filepath" "strings" + + "github.com/huly-stream/internal/pkg/resconv" ) -// HLSManifest represents an HLS manifest file -// with metadata about the playlist and its segments. -type HLSManifest struct { - Version int - TargetDuration int - SequenceNumber int - Segments []Segment - EndList bool -} +// GenerateHLSPlaylist generates master file for master files for resolution levels +func GenerateHLSPlaylist(levels []string, outputPath, uploadID string) error { + p := filepath.Join(outputPath, uploadID, fmt.Sprintf("%v_master.m3u8", uploadID)) + d := filepath.Dir(p) + _ = os.MkdirAll(d, os.ModePerm) + // #nosec + file, err := os.Create(p) + if err != nil { + return err + } + defer func() { _ = file.Close() }() -// Segment represents a media segment in the HLS manifest. -type Segment struct { - URI string - Duration float64 - Title string -} - -// ToM3U8 serializes the HLSManifest to an M3U8 file format. -func (m *HLSManifest) ToM3U8() string { - var builder strings.Builder - - builder.WriteString("#EXTM3U\n") - builder.WriteString(fmt.Sprintf("#EXT-X-VERSION:%d\n", m.Version)) - builder.WriteString(fmt.Sprintf("#EXT-X-TARGETDURATION:%d\n", m.TargetDuration)) - builder.WriteString(fmt.Sprintf("#EXT-X-MEDIA-SEQUENCE:%d\n", m.SequenceNumber)) - - for _, segment := range m.Segments { - if segment.Title != "" { - builder.WriteString(fmt.Sprintf("#EXTINF:%.2f,%s\n", segment.Duration, segment.Title)) - } else { - builder.WriteString(fmt.Sprintf("#EXTINF:%.2f,\n", segment.Duration)) - } - builder.WriteString(fmt.Sprintf("%s\n", segment.URI)) + _, err = file.WriteString("#EXTM3U\n") + if err != nil { + return err } - if m.EndList { - builder.WriteString("#EXT-X-ENDLIST\n") - } + for _, res := range levels { + var bandwidth = resconv.Bandwidth(res) + var resolution = strings.ReplaceAll(resconv.Resolution(res), ":", "x") - return builder.String() -} - -// FromM3U8 converts raw input to the hls master file -// nolint -func FromM3U8(data string) (*HLSManifest, error) { - scanner := bufio.NewScanner(strings.NewReader(data)) - manifest := &HLSManifest{} - var currentSegment *Segment - - for scanner.Scan() { - line := strings.TrimSpace(scanner.Text()) - if line == "" { - continue + _, err = file.WriteString(fmt.Sprintf("#EXT-X-STREAM-INF:BANDWIDTH=%d,RESOLUTION=%v\n", bandwidth, resolution)) + if err != nil { + return err } - if strings.HasPrefix(line, "#EXTM3U") { - continue - } - if strings.HasPrefix(line, "#EXT-X-VERSION:") { - version, err := strconv.Atoi(strings.TrimPrefix(line, "#EXT-X-VERSION:")) - if err != nil { - return nil, err - } - manifest.Version = version - } else if strings.HasPrefix(line, "#EXT-X-TARGETDURATION:") { - targetDuration, err := strconv.Atoi(strings.TrimPrefix(line, "#EXT-X-TARGETDURATION:")) - if err != nil { - return nil, err - } - manifest.TargetDuration = targetDuration - } else if strings.HasPrefix(line, "#EXT-X-MEDIA-SEQUENCE:") { - sequenceNumber, err := strconv.Atoi(strings.TrimPrefix(line, "#EXT-X-MEDIA-SEQUENCE:")) - if err != nil { - return nil, err - } - manifest.SequenceNumber = sequenceNumber - } else if strings.HasPrefix(line, "#EXTINF:") { - parts := strings.SplitN(strings.TrimPrefix(line, "#EXTINF:"), ",", 2) - duration, err := strconv.ParseFloat(parts[0], 64) - if err != nil { - return nil, err - } - title := "" - if len(parts) > 1 { - title = parts[1] - } - currentSegment = &Segment{Duration: duration, Title: title} - } else if strings.HasPrefix(line, "#EXT-X-ENDLIST") { - manifest.EndList = true - } else if currentSegment != nil { - currentSegment.URI = line - manifest.Segments = append(manifest.Segments, *currentSegment) - currentSegment = nil + _, err = file.WriteString(fmt.Sprintf("%s_%s_master.m3u8\n", uploadID, res)) + if err != nil { + return err } } - if err := scanner.Err(); err != nil { - return nil, err - } - - return manifest, nil + return nil } diff --git a/internal/pkg/manifest/hls_test.go b/internal/pkg/manifest/hls_test.go index 15f3dbe022..624e3a27a9 100644 --- a/internal/pkg/manifest/hls_test.go +++ b/internal/pkg/manifest/hls_test.go @@ -14,130 +14,38 @@ package manifest_test import ( + "os" + "path/filepath" "testing" "github.com/huly-stream/internal/pkg/manifest" - "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) -func TestToM3U8(t *testing.T) { - tests := []struct { - name string - manifest manifest.HLSManifest - expected string - }{ - { - name: "simple manifest", - manifest: manifest.HLSManifest{ - Version: 3, - TargetDuration: 10, - SequenceNumber: 1, - Segments: []manifest.Segment{ - {URI: "segment1.ts", Duration: 9.5, Title: "Segment 1"}, - {URI: "segment2.ts", Duration: 9.0, Title: "Segment 2"}, - }, - EndList: true, - }, - expected: `#EXTM3U -#EXT-X-VERSION:3 -#EXT-X-TARGETDURATION:10 -#EXT-X-MEDIA-SEQUENCE:1 -#EXTINF:9.50,Segment 1 -segment1.ts -#EXTINF:9.00,Segment 2 -segment2.ts -#EXT-X-ENDLIST -`, - }, - { - name: "empty manifest", - manifest: manifest.HLSManifest{ - Version: 3, - TargetDuration: 10, - SequenceNumber: 1, - Segments: []manifest.Segment{}, - EndList: false, - }, - expected: `#EXTM3U -#EXT-X-VERSION:3 -#EXT-X-TARGETDURATION:10 -#EXT-X-MEDIA-SEQUENCE:1 -`, - }, +func TestGenerateHLSPlaylist(t *testing.T) { + resolutions := []string{"320p", "480p", "720p", "1080p", "4k", "8k"} + uploadID := "test123" + + err := manifest.GenerateHLSPlaylist(resolutions, "", uploadID) + require.NoError(t, err) + + outputPath := filepath.Join(uploadID, uploadID+"_master.m3u8") + + _, err = os.Stat(outputPath) + require.NoError(t, err, "Master playlist file should exist") + + // #nosec + data, err := os.ReadFile(outputPath) + require.NoError(t, err, "Error reading the generated file") + + playlistContent := string(data) + + require.Contains(t, playlistContent, "#EXTM3U", "File must start with #EXTM3U") + + for _, res := range resolutions { + expectedLine := uploadID + "_" + res + "_master.m3u8" + require.Contains(t, playlistContent, expectedLine, "Missing expected reference: "+expectedLine) } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - actual := tt.manifest.ToM3U8() - assert.Equal(t, tt.expected, actual) - }) - } -} - -func TestFromM3U8(t *testing.T) { - tests := []struct { - name string - data string - expected manifest.HLSManifest - err bool - }{ - { - name: "valid manifest", - data: `#EXTM3U -#EXT-X-VERSION:3 -#EXT-X-TARGETDURATION:10 -#EXT-X-MEDIA-SEQUENCE:1 -#EXTINF:9.50,Segment 1 -segment1.ts -#EXTINF:9.00,Segment 2 -segment2.ts -#EXT-X-ENDLIST -`, - expected: manifest.HLSManifest{ - Version: 3, - TargetDuration: 10, - SequenceNumber: 1, - Segments: []manifest.Segment{ - {URI: "segment1.ts", Duration: 9.5, Title: "Segment 1"}, - {URI: "segment2.ts", Duration: 9.0, Title: "Segment 2"}, - }, - EndList: true, - }, - err: false, - }, - // { - // name: "missing target duration", - // data: `#EXTM3U - // #EXT-X-VERSION:3 - // #EXT-X-MEDIA-SEQUENCE:1 - // #EXTINF:9.50,Segment 1 - // segment1.ts - // `, - // expected: manifest.HLSManifest{}, - // err: true, - // }, - { - name: "empty file", - data: "", - expected: manifest.HLSManifest{ - Version: 0, - TargetDuration: 0, - SequenceNumber: 0, - EndList: false, - }, - err: false, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - actual, err := manifest.FromM3U8(tt.data) - if tt.err { - assert.Error(t, err) - } else { - assert.NoError(t, err) - assert.Equal(t, tt.expected, *actual) - } - }) - } + _ = os.RemoveAll(uploadID) } diff --git a/internal/pkg/resconv/resconv.go b/internal/pkg/resconv/resconv.go new file mode 100644 index 0000000000..4e33a8db44 --- /dev/null +++ b/internal/pkg/resconv/resconv.go @@ -0,0 +1,129 @@ +// +// Copyright © 2025 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +// Package resconv implements conversions to and from string representations of video resolutions. +package resconv + +import ( + "sort" + "strconv" + "strings" +) + +const defaultLevel = "320p" + +var prefixes = []struct { + pixels int + label string +}{ + {pixels: 640 * 480, label: "320p"}, + {pixels: 1280 * 720, label: "480p"}, + {pixels: 1920 * 1080, label: "720p"}, + {pixels: 2560 * 1440, label: "1080p"}, + {pixels: 3840 * 2160, label: "2k"}, + {pixels: 5120 * 2880, label: "4k"}, + {pixels: 7680 * 4320, label: "5k"}, +} + +var bandwidthMap = map[string]int{ + "320p": 300000, + "360p": 500000, + "480p": 2000000, + "720p": 5000000, + "1080p": 8000000, + "1440p": 16000000, + "4k": 25000000, + "8k": 50000000, +} + +var resolutions = map[string]string{ + "320p": "480:240", + "480p": "640:480", + "720p": "1280:720", + "1080p": "1920:1080", + "2k": "2048:1080", + "4k": "3840:2160", + "5k": "5120:2880", + "8k": "7680:4320", +} + +// SubLevels returns sublevels for the resolution +func SubLevels(resolution string) (res []string) { + var pixels = Pixels(resolution) + var idx = sort.Search(len(prefixes), func(i int) bool { + return pixels < prefixes[i].pixels + }) + if idx < 2 { + return res + } + + idx-- + idx = min(idx, 3) + + for idx >= 1 { + res = append(res, prefixes[idx].label) + idx-- + if len(res) == 2 { + break + } + } + + return res +} + +// Resolution returns default resolution based on the level +func Resolution(level string) string { + if v, ok := resolutions[level]; ok { + return v + } + return Resolution(defaultLevel) +} + +// Level converts the resolution to short prefix +func Level(resolution string) string { + var pixels = Pixels(resolution) + idx := sort.Search(len(prefixes), func(i int) bool { + return pixels < prefixes[i].pixels + }) + if idx == len(prefixes) { + return "8k" + } + + return prefixes[idx].label +} + +// Pixels returns amount of pixels for the resolution +func Pixels(resolution string) int { + var parts = strings.Split(resolution, ":") + var w, h = 420, 240 + + if len(parts) > 1 { + var _w, _ = strconv.Atoi(parts[0]) + var _h, _ = strconv.Atoi(parts[1]) + w = max(w, _w) + h = max(h, _h) + } + + return w * h +} + +// Bandwidth returns default bandwidth for the resolution +func Bandwidth(resolution string) int { + if v, ok := bandwidthMap[resolution]; ok { + return v + } + + return bandwidthMap[defaultLevel] +} diff --git a/internal/pkg/resconv/resconv_test.go b/internal/pkg/resconv/resconv_test.go new file mode 100644 index 0000000000..4139e21370 --- /dev/null +++ b/internal/pkg/resconv/resconv_test.go @@ -0,0 +1,120 @@ +// Copyright © 2025 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. + +package resconv_test + +import ( + "testing" + + "github.com/huly-stream/internal/pkg/resconv" + "github.com/stretchr/testify/require" +) + +func Test_Resconv_ShouldReturnCorrectPrefix(t *testing.T) { + tests := []struct { + res string + expected string + }{ + {res: "320:240", expected: "320p"}, + {res: "640:480", expected: "480p"}, + {res: "1280:720", expected: "720p"}, + {res: "1920:1080", expected: "1080p"}, + {res: "2560:1440", expected: "2k"}, + {res: "3840:2160", expected: "4k"}, + {res: "5120:2880", expected: "5k"}, + {res: "9000:4000", expected: "8k"}, + } + + for _, tt := range tests { + t.Run(tt.expected, func(t *testing.T) { + result := resconv.Level(tt.res) + require.Equal(t, tt.expected, result, "ResolutionFromPixels(%d)", tt.res) + }) + } +} + +func Test_Resconv_ShouldReturnCorrectPrefixes(t *testing.T) { + tests := []struct { + name string + res string + expected []string + }{ + { + name: "pixels below smallest resolution", + res: "640:479", + expected: nil, + }, + { + name: "pixels equal to smallest resolution", + res: "640:480", + expected: nil, + }, + { + name: "pixels just above smallest resolution", + res: "641:480", + expected: nil, + }, + { + name: "pixels equal to 720p", + res: "1280:720", + expected: []string{"480p"}, + }, + { + name: "pixels just above 480p", + res: "1280:721", + expected: []string{"480p"}, + }, + { + name: "pixels equal to 1k", + res: "1920:1080", + expected: []string{"720p", "480p"}, + }, + { + name: "pixels just above 1k", + res: "1920:1081", + expected: []string{"720p", "480p"}, + }, + { + name: "pixels equal to 1k", + res: "2560:1440", + expected: []string{"1080p", "720p"}, + }, + { + name: "pixels equal to 2k", + res: "3840:2160", + expected: []string{"1080p", "720p"}, + }, + { + name: "pixels equal to 4k", + res: "5120:2160", + expected: []string{"1080p", "720p"}, + }, + { + name: "pixels equal to 5k", + res: "7680:4320", + expected: []string{"1080p", "720p"}, + }, + { + name: "pixels above largest resolution", + res: "7681:4320", + expected: []string{"1080p", "720p"}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + result := resconv.SubLevels(tt.res) + require.Equal(t, tt.expected, result, "SubResolutionsFromPixels(%d) returned unexpected result", tt.res) + }) + } +} diff --git a/internal/pkg/transcoding/command.go b/internal/pkg/transcoding/command.go index f43b5cf594..0065fa3196 100644 --- a/internal/pkg/transcoding/command.go +++ b/internal/pkg/transcoding/command.go @@ -22,61 +22,29 @@ import ( "os" "os/exec" "path/filepath" - "sort" - "strconv" - "strings" "github.com/pkg/errors" "github.com/huly-stream/internal/pkg/log" + "github.com/huly-stream/internal/pkg/resconv" "go.uber.org/zap" ) // Options represents configuration for the ffmpeg command type Options struct { - OuputDir string - Resolutions []string - Threads int - UploadID string + OuputDir string + ScalingLevels []string + Level string + Threads int + UploadID string } -func measure(options *Options) int64 { - var res int64 - for _, resolution := range options.Resolutions { - var w, h int - var parts = strings.Split(resolution, ":") - - if len(parts) > 1 { - w, _ = strconv.Atoi(parts[0]) - w = max(w, 320) - h, _ = strconv.Atoi(parts[1]) - h = max(h, 240) - - res += int64(w) * int64(h) - } - } - - return max(res, 320*240) -} - -func newFfmpegCommand(ctx context.Context, in io.Reader, options *Options) (*exec.Cmd, error) { - if options == nil { - return nil, errors.New("options should not be nil") - } +func newFfmpegCommand(ctx context.Context, in io.Reader, args []string) (*exec.Cmd, error) { if ctx == nil { return nil, errors.New("ctx should not be nil") } - var logger = log.FromContext(ctx).With(zap.String("func", "NewFFMpegCommand")) - var args []string - - if options.Resolutions == nil { - logger.Debug("resolutions were not provided, building audio command...") - args = BuildAudioCommand(options) - } else { - logger.Debug("building video command...") - args = BuildVideoCommand(options) - } + var logger = log.FromContext(ctx).With(zap.String("func", "newFFMpegCommand")) logger.Debug("prepared command: ", zap.Strings("args", args)) @@ -90,6 +58,7 @@ func newFfmpegCommand(ctx context.Context, in io.Reader, options *Options) (*exe func buildCommonComamnd(opts *Options) []string { return []string{ + "-nostdin", "-threads", fmt.Sprint(opts.Threads), "-i", "pipe:0", } @@ -105,25 +74,25 @@ func BuildAudioCommand(opts *Options) []string { ) } -// BuildVideoCommand returns flags for ffmpeg for video transcoding -func BuildVideoCommand(opts *Options) []string { +// BuildRawVideoCommand returns an extremely lightweight ffmpeg command for converting raw video without extra cost. +func BuildRawVideoCommand(opts *Options) []string { + return append(buildCommonComamnd(opts), + "-c:v", + "copy", + "-fps_mode", + "vfr", + "-hls_time", "5", + "-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))) +} + +// BuildScalingVideoCommand returns flags for ffmpeg for video scaling +func BuildScalingVideoCommand(opts *Options) []string { var result = buildCommonComamnd(opts) - - for _, res := range opts.Resolutions { - var prefix string - var w, h int - var parts = strings.Split(res, ":") - - if len(parts) > 1 { - w, _ = strconv.Atoi(parts[0]) - h, _ = strconv.Atoi(parts[1]) - } - w = max(w, 640) - h = max(h, 480) - prefix = ResolutionFromPixels(w * h) - + for _, level := range opts.ScalingLevels { result = append(result, - "-vf", fmt.Sprintf("scale=%d:%d", w, h), + "-vf", "scale="+resconv.Resolution(level), "-c:v", "libx264", "-preset", "veryfast", @@ -131,33 +100,9 @@ func BuildVideoCommand(opts *Options) []string { "-g", "60", "-hls_time", "5", "-hls_list_size", "0", - "-hls_segment_filename", filepath.Join(opts.OuputDir, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", prefix)), - filepath.Join(opts.OuputDir, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, prefix))) + "-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))) } return result } - -var resolutions = []struct { - pixels int - label string -}{ - {pixels: 640 * 480, label: "320p"}, - {pixels: 1280 * 720, label: "480p"}, - {pixels: 1920 * 1080, label: "720p"}, - {pixels: 2560 * 1440, label: "1k"}, - {pixels: 3840 * 2160, label: "2k"}, - {pixels: 5120 * 2160, label: "4k"}, - {pixels: 7680 * 4320, label: "5k"}, -} - -// ResolutionFromPixels converts pixel count to short string -func ResolutionFromPixels(pixels int) string { - idx := sort.Search(len(resolutions), func(i int) bool { - return pixels < resolutions[i].pixels - }) - if idx == len(resolutions) { - return "8k" - } - return resolutions[idx].label -} diff --git a/internal/pkg/transcoding/command_test.go b/internal/pkg/transcoding/command_test.go index 14a7db515d..719c38467a 100644 --- a/internal/pkg/transcoding/command_test.go +++ b/internal/pkg/transcoding/command_test.go @@ -14,49 +14,36 @@ package transcoding_test import ( - "runtime" "strings" "testing" + "github.com/huly-stream/internal/pkg/resconv" "github.com/huly-stream/internal/pkg/transcoding" "github.com/stretchr/testify/require" ) -func Test_BuildVideoCommand_Basic(t *testing.T) { - if runtime.GOOS == "windows" { - t.Skip() - } - var simpleHlsCommand = transcoding.BuildVideoCommand(&transcoding.Options{ - OuputDir: "test", - UploadID: "1", - Threads: 4, - Resolutions: []string{"1280:720"}, +func Test_BuildVideoCommand_Scaling(t *testing.T) { + var scaleCommand = transcoding.BuildScalingVideoCommand(&transcoding.Options{ + OuputDir: "test", + UploadID: "1", + Threads: 4, + ScalingLevels: []string{"720p", "480p"}, }) - const expected = `-threads 4 -i pipe:0 -vf scale=1280:720 -c:v libx264 -preset veryfast -crf 23 -g 60 -hls_time 5 -hls_list_size 0 -hls_segment_filename test/1_%03d_720p.ts test/1_720p_master.m3u8` + const expected = `-nostdin -threads 4 -i pipe:0 -vf scale=1280:720 -c:v libx264 -preset veryfast -crf 23 -g 60 -hls_time 5 -hls_list_size 0 -hls_segment_filename test/1/1_%03d_720p.ts test/1/1_720p_master.m3u8 -vf scale=640:480 -c:v libx264 -preset veryfast -crf 23 -g 60 -hls_time 5 -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(simpleHlsCommand, " ")) + require.Contains(t, expected, strings.Join(scaleCommand, " ")) } -func TestResolutionFromPixels(t *testing.T) { - tests := []struct { - pixels int - expected string - }{ - {pixels: 320 * 240, expected: "320p"}, - {pixels: 640 * 480, expected: "480p"}, - {pixels: 1280 * 720, expected: "720p"}, - {pixels: 1920 * 1080, expected: "1k"}, - {pixels: 2560 * 1440, expected: "2k"}, - {pixels: 3840 * 2160, expected: "4k"}, - {pixels: 5120 * 2160, expected: "5k"}, - {pixels: 9000 * 4000, expected: "8k"}, - } +func Test_BuildVideoCommand_Raw(t *testing.T) { + var rawCommand = transcoding.BuildRawVideoCommand(&transcoding.Options{ + OuputDir: "test", + UploadID: "1", + Threads: 4, + Level: resconv.Level("651:490"), + }) - for _, tt := range tests { - t.Run(tt.expected, func(t *testing.T) { - result := transcoding.ResolutionFromPixels(tt.pixels) - require.Equal(t, tt.expected, result, "ResolutionFromPixels(%d)", tt.pixels) - }) - } + const expected = `-nostdin -threads 4 -i pipe:0 -c:v copy -fps_mode vfr -hls_time 5 -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/transcoding/scheduler.go b/internal/pkg/transcoding/scheduler.go index 25872a3047..080aeedbb5 100644 --- a/internal/pkg/transcoding/scheduler.go +++ b/internal/pkg/transcoding/scheduler.go @@ -17,14 +17,15 @@ package transcoding import ( "context" - "strings" "sync" + "time" "github.com/pkg/errors" "github.com/google/uuid" "github.com/huly-stream/internal/pkg/config" "github.com/huly-stream/internal/pkg/log" + "github.com/huly-stream/internal/pkg/resconv" "github.com/huly-stream/internal/pkg/sharedpipe" "github.com/huly-stream/internal/pkg/uploader" "github.com/tus/tusd/v2/pkg/handler" @@ -40,6 +41,7 @@ type Scheduler struct { mainContext context.Context logger *zap.Logger workers sync.Map + cancels sync.Map } // NewScheduler creates a new scheduler for transcode operations. @@ -57,58 +59,95 @@ func (s *Scheduler) NewUpload(ctx context.Context, info handler.FileInfo) (handl if info.ID == "" { info.ID = uuid.NewString() } - + s.logger.Sugar().Debugf("upload: %v", info) s.logger.Debug("NewUpload", zap.String("ID", info.ID)) - var result = &Worker{ - done: make(chan struct{}), + var worker = &Worker{ writer: sharedpipe.NewWriter(), info: info, - logger: log.FromContext(s.mainContext).With(zap.String("Worker", info.ID)), + logger: log.FromContext(s.mainContext).With(zap.String("worker", info.ID)), + done: make(chan struct{}), } - var resolutions = strings.Split(info.MetaData["resolutions"], ",") + var scaling = resconv.SubLevels(info.MetaData["resolution"]) + var level = resconv.Level(info.MetaData["resolution"]) + var cost int64 + + for _, scale := range scaling { + cost += int64(resconv.Pixels(resconv.Resolution(scale))) + } + + if !s.limiter.TryConsume(cost) { + s.logger.Debug("run out of resources for scaling") + scaling = nil + } var commandOptions = Options{ - OuputDir: s.conf.OutputDir, - Threads: s.conf.MaxThreads, - UploadID: info.ID, - Resolutions: resolutions, - } - - result.cost = measure(&commandOptions) - - if !s.limiter.TryConsume(result.cost) { - s.logger.Error("run out of resources") - return nil, errors.New("run out of resources") + OuputDir: s.conf.OutputDir, + Threads: s.conf.MaxThreads, + UploadID: info.ID, + Level: level, + ScalingLevels: scaling, } if s.conf.EndpointURL != nil { - s.logger.Debug("found endpoint url in the config, starting uploader...") - var contentUploader, err = uploader.New(s.mainContext, *s.conf, info.ID, info.MetaData) + s.logger.Sugar().Debugf("initializing uploader for %v", info) + var contentUploader, err = uploader.New(s.mainContext, s.conf.OutputDir, s.conf.EndpointURL, info) if err != nil { + s.logger.Error("can not create uploader", zap.Error(err)) return nil, err } - result.contentUploader = contentUploader + + worker.contentUploader = contentUploader go func() { - var serverErr = result.contentUploader.Serve() - result.logger.Debug("content uploader has finished", zap.Error(serverErr)) + var serverErr = worker.contentUploader.Serve() + worker.logger.Debug("content uploader has finished", zap.Error(serverErr)) }() } - - s.workers.Store(result.info.ID, result) - s.logger.Sugar().Debugf("New Upload: info %v", result.info) - if err := result.start(s.mainContext, &commandOptions); err != nil { + s.workers.Store(worker.info.ID, worker) + if err := worker.start(s.mainContext, &commandOptions); err != nil { return nil, err } - return result, nil + + go func() { + worker.wg.Wait() + s.limiter.ReturnCapacity(cost) + s.logger.Debug("returned capacity", zap.Int64("capacity", cost)) + close(worker.done) + }() + + s.logger.Debug("NewUpload", zap.String("done", info.ID)) + return worker, nil } // GetUpload returns current a worker based on upload id func (s *Scheduler) GetUpload(ctx context.Context, id string) (upload handler.Upload, err error) { if v, ok := s.workers.Load(id); ok { s.logger.Debug("GetUpload: found worker by id", zap.String("id", id)) - return v.(*Worker), nil + var w = v.(*Worker) + var cancelCtx, cancel = context.WithCancel(context.Background()) + if v, ok := s.cancels.Load(id); ok { + v.(context.CancelFunc)() + } + s.cancels.Store(id, cancel) + go func() { + select { + case <-w.done: + w.logger.Debug("upload timeout just canceled") + s.cancels.Delete(id) + return + case <-cancelCtx.Done(): + w.logger.Debug("upload refreshed") + return + case <-time.After(s.conf.Timeout): + w.logger.Debug("upload timeout") + s.cancels.Delete(id) + var terminateCtx, terminateCancel = context.WithTimeout(context.Background(), s.conf.Timeout) + defer terminateCancel() + _ = w.Terminate(terminateCtx) + } + }() + return w, nil } s.logger.Debug("GetUpload: worker not found", zap.String("id", id)) return nil, errors.New("bad id") @@ -117,8 +156,7 @@ func (s *Scheduler) GetUpload(ctx context.Context, id string) (upload handler.Up // AsTerminatableUpload returns tusd handler.TerminatableUpload func (s *Scheduler) AsTerminatableUpload(upload handler.Upload) handler.TerminatableUpload { var worker = upload.(*Worker) - s.logger.Debug("AsTerminatableUpload, trying to return capacity", zap.Int64("cost", worker.cost)) - s.limiter.ReturnCapacity(worker.cost) + s.logger.Debug("AsTerminatableUpload") return worker } diff --git a/internal/pkg/transcoding/worker.go b/internal/pkg/transcoding/worker.go index 60baccc465..1235bfcf02 100644 --- a/internal/pkg/transcoding/worker.go +++ b/internal/pkg/transcoding/worker.go @@ -17,9 +17,11 @@ package transcoding import ( "context" "io" + "sync" "github.com/pkg/errors" + "github.com/huly-stream/internal/pkg/manifest" "github.com/huly-stream/internal/pkg/sharedpipe" "github.com/huly-stream/internal/pkg/uploader" "github.com/tus/tusd/v2/pkg/handler" @@ -33,8 +35,9 @@ type Worker struct { info handler.FileInfo writer *sharedpipe.Writer reader *sharedpipe.Reader - cost int64 - done chan struct{} + + wg sync.WaitGroup + done chan struct{} } // WriteChunk calls when client sends a chunk of raw data @@ -73,7 +76,7 @@ func (w *Worker) Terminate(ctx context.Context) error { w.logger.Debug("Terminating...") if w.contentUploader != nil { go func() { - <-w.done + w.wg.Wait() w.contentUploader.Rollback() }() } @@ -94,7 +97,7 @@ func (w *Worker) FinishUpload(ctx context.Context) error { w.logger.Debug("finishing upload...") if w.contentUploader != nil { go func() { - <-w.done + w.wg.Wait() w.contentUploader.Terminate() }() } @@ -108,18 +111,49 @@ func (s *Scheduler) AsConcatableUpload(upload handler.Upload) handler.Concatable } func (w *Worker) start(ctx context.Context, options *Options) error { + defer w.logger.Debug("start done") w.reader = w.writer.Transpile() - var cmd, err = newFfmpegCommand(ctx, w.reader, options) - if err != nil { + + if err := manifest.GenerateHLSPlaylist(append(options.ScalingLevels, options.Level), options.OuputDir, options.UploadID); err != nil { return err } + + w.wg.Add(1) go func() { - defer close(w.done) - if runErr := cmd.Run(); runErr != nil { - w.logger.Error("transoding provider is exited with error", zap.Error(err)) - } else { - w.logger.Debug("transoding provider has finished without errors") + defer w.wg.Done() + var logger = w.logger.With(zap.String("command", "raw")) + defer logger.Debug("done") + + var args = BuildRawVideoCommand(options) + var convertSourceCommand, err = newFfmpegCommand(ctx, w.reader, args) + if err != nil { + logger.Debug("can not start", zap.Error(err)) + } + err = convertSourceCommand.Run() + if err != nil { + logger.Debug("finished with error", zap.Error(err)) } }() + + if len(options.ScalingLevels) > 0 { + w.wg.Add(1) + var scalingCommandReader = w.writer.Transpile() + go func() { + defer w.wg.Done() + var logger = w.logger.With(zap.String("command", "scaling")) + defer logger.Debug("done") + + var args = BuildScalingVideoCommand(options) + var convertSourceCommand, err = newFfmpegCommand(ctx, scalingCommandReader, args) + if err != nil { + logger.Debug("can not start", zap.Error(err)) + } + err = convertSourceCommand.Run() + if err != nil { + logger.Debug("finished with error", zap.Error(err)) + } + }() + } + return nil } diff --git a/internal/pkg/uploader/datalake.go b/internal/pkg/uploader/datalake.go index 43ba535efd..dfff697012 100644 --- a/internal/pkg/uploader/datalake.go +++ b/internal/pkg/uploader/datalake.go @@ -18,7 +18,6 @@ import ( "context" "io" "mime/multipart" - "net/url" "os" "path/filepath" @@ -116,6 +115,7 @@ func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error client := fasthttp.Client{} if err := client.Do(req, res); err != nil { + logger.Error("failed to del", zap.Error(err)) return errors.Wrapf(err, "delete failed") } @@ -126,5 +126,5 @@ func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error func getObjectKey(s string) string { var _, objectKey = filepath.Split(s) - return url.QueryEscape(objectKey) + return objectKey } diff --git a/internal/pkg/uploader/s3.go b/internal/pkg/uploader/s3.go index ba688618f9..faf15e22d0 100644 --- a/internal/pkg/uploader/s3.go +++ b/internal/pkg/uploader/s3.go @@ -37,6 +37,7 @@ import ( type S3Storage struct { client *s3.Client bucketName string + logger *zap.Logger } // NewS3 creates a new S3 storage @@ -44,6 +45,7 @@ func NewS3(ctx context.Context, endpoint string) Storage { var accessKeyID = os.Getenv("AWS_ACCESS_KEY_ID") var accessKeySecret = os.Getenv("AWS_SECRET_ACCESS_KEY") var bucketName = os.Getenv("AWS_BUCKET_NAME") + var logger = log.FromContext(ctx).With(zap.String("s3", "storage")) cfg, err := config.LoadDefaultConfig(ctx, config.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(accessKeyID, accessKeySecret, "")), @@ -61,6 +63,7 @@ func NewS3(ctx context.Context, endpoint string) Storage { return &S3Storage{ client: s3Client, bucketName: bucketName, + logger: logger, } } diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index 86735191a4..f61a6396ee 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -18,52 +18,54 @@ import ( "context" "net/url" "os" + "path/filepath" "strings" "sync" "time" "github.com/pkg/errors" + "github.com/tus/tusd/v2/pkg/handler" "github.com/fsnotify/fsnotify" - "github.com/huly-stream/internal/pkg/config" "github.com/huly-stream/internal/pkg/log" "go.uber.org/zap" ) type uploader struct { - ctx context.Context - cancel context.CancelFunc - baseDir string - uploadID string - masterFiles sync.Map - postponeDuration time.Duration - sentFiles sync.Map - storage Storage - contexts sync.Map - retryCount int - removeLocalContentOnUpload bool - eventBufferCount uint - isMasterFileFunc func(s string) bool + ctx context.Context + cancel context.CancelFunc + baseDir string + uploadID string + masterFiles sync.Map + postponeDuration time.Duration + sentFiles sync.Map + storage Storage + contexts sync.Map + retryCount int + eventBufferCount uint + isMasterFileFunc func(s string) bool +} + +func (u *uploader) retry(action func() error) { + for range u.retryCount { + if err := action(); err == nil { + return + } + } } // Rollback deletes all delivered files and also deletes all local content by uploadID func (u *uploader) Rollback() { - log.FromContext(u.ctx).Debug("cancel") + logger := log.FromContext(u.ctx).With(zap.String("uploader", "Rollback")) + logger.Debug("starting") defer u.cancel() + u.sentFiles.Range(func(key, value any) bool { - log.FromContext(u.ctx).Debug("deleting remote file", zap.String("key", key.(string))) - for range u.retryCount { - var err = u.storage.DeleteFile(u.ctx, key.(string)) - if err == nil { - break - } - log.FromContext(u.ctx).Debug("can not delete file", zap.Error(err)) - } + logger.Debug("deleting remote file", zap.String("key", key.(string))) + u.retry(func() error { return u.storage.DeleteFile(u.ctx, key.(string)) }) return true }) - if !u.removeLocalContentOnUpload { - return - } + u.sentFiles.Range(func(key, value any) bool { log.FromContext(u.ctx).Debug("deleting local file", zap.String("key", key.(string))) _ = os.Remove(key.(string)) @@ -72,50 +74,47 @@ func (u *uploader) Rollback() { } func (u *uploader) Terminate() { - log.FromContext(u.ctx).Debug("terminate") + logger := log.FromContext(u.ctx).With(zap.String("uploader", "Terminate")) + logger.Debug("starting") defer u.cancel() + u.masterFiles.Range(func(key, value any) bool { log.FromContext(u.ctx).Debug("uploading master file", zap.String("key", key.(string))) - for range u.retryCount { - var uploadErr = u.storage.UploadFile(u.ctx, key.(string)) - if uploadErr == nil { - break - } - log.FromContext(u.ctx).Debug("can not upload file", zap.Error(uploadErr)) - } + go u.retry(func() error { return u.storage.UploadFile(u.ctx, key.(string)) }) return true }) - if !u.removeLocalContentOnUpload { - return - } + u.masterFiles.Range(func(key, value any) bool { - log.FromContext(u.ctx).Debug("deleting local master file", zap.String("key", key.(string))) _ = os.Remove(key.(string)) return true }) + u.sentFiles.Range(func(key, value any) bool { - log.FromContext(u.ctx).Debug("deleting local file", zap.String("key", key.(string))) _ = os.Remove(key.(string)) return true }) } func (u *uploader) Serve() error { - var logger = log.FromContext(u.ctx) - logger = logger.With(zap.String("uploader", u.uploadID), zap.String("dir", u.baseDir)) + var logger = log.FromContext(u.ctx).With(zap.String("uploader", u.uploadID), zap.String("dir", u.baseDir)) var watcher, err = fsnotify.NewBufferedWatcher(u.eventBufferCount) + if err != nil { logger.Error("can not start watcher") return err } + + _ = os.MkdirAll(u.baseDir, os.ModePerm) + if err := watcher.Add(u.baseDir); err != nil { return err } + defer func() { _ = watcher.Close() }() - logger.Debug("uploader initialized and started to watch") + logger.Debug("the uploader has initialized and started watching") for { select { @@ -123,6 +122,9 @@ func (u *uploader) Serve() error { logger.Debug("done") return u.ctx.Err() case event, ok := <-watcher.Events: + if strings.HasSuffix(event.Name, "tmp") { + continue + } if !strings.Contains(event.Name, u.uploadID) { continue } @@ -169,29 +171,32 @@ type Storage interface { } // New creates a new instance of Uplaoder -func New(ctx context.Context, conf config.Config, uploadID string, metadata map[string]string) (Uploader, error) { - var uploaderCtx, uploaderCancel = context.WithCancel(ctx) +func New(ctx context.Context, baseDir string, endpointURL *url.URL, uploadInfo handler.FileInfo) (Uploader, error) { + var uploaderCtx, uploadCancel = context.WithCancel(context.Background()) + uploaderCtx = log.WithLoggerFields(uploaderCtx) + go func() { + <-ctx.Done() + time.Sleep(time.Minute * 2) + uploadCancel() + }() var storage Storage var err error - if conf.EndpointURL != nil { - storage, err = NewStorageByURL(ctx, conf.EndpointURL, metadata) - if err != nil { - uploaderCancel() - return nil, err - } + storage, err = NewStorageByURL(ctx, endpointURL, uploadInfo.MetaData) + if err != nil { + uploadCancel() + return nil, err } return &uploader{ - ctx: uploaderCtx, - cancel: uploaderCancel, - uploadID: uploadID, - removeLocalContentOnUpload: conf.RemoveContentOnUpload, - postponeDuration: time.Second * 2, - storage: storage, - retryCount: 5, - baseDir: conf.OutputDir, - eventBufferCount: 100, + ctx: uploaderCtx, + cancel: uploadCancel, + uploadID: uploadInfo.ID, + postponeDuration: time.Second * 2, + storage: storage, + retryCount: 5, + baseDir: filepath.Join(baseDir, uploadInfo.ID), + eventBufferCount: 100, isMasterFileFunc: func(s string) bool { return strings.HasSuffix(s, "m3u8") },