mirror of
https://github.com/hcengineering/platform.git
synced 2026-10-01 05:55:09 +02:00
Open telemetry support
Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>
This commit is contained in:
+13
-19
@@ -16,6 +16,7 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"sync"
|
||||
|
||||
@@ -24,8 +25,6 @@ import (
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/getsentry/sentry-go"
|
||||
sentryhttp "github.com/getsentry/sentry-go/http"
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/hcengineering/stream/internal/pkg/api/v1/recording"
|
||||
@@ -33,6 +32,7 @@ import (
|
||||
"github.com/hcengineering/stream/internal/pkg/config"
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/queue"
|
||||
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
|
||||
)
|
||||
|
||||
func main() {
|
||||
@@ -53,16 +53,14 @@ func main() {
|
||||
}
|
||||
logger.Sugar().Debug("using config", zap.Any("config", cfg))
|
||||
|
||||
if cfg.SentryDsn != "" {
|
||||
if err := sentry.Init(sentry.ClientOptions{
|
||||
Dsn: cfg.SentryDsn,
|
||||
Tags: map[string]string{"application": "stream"},
|
||||
}); err != nil {
|
||||
logger.Sugar().Fatalf("sentry.Init: %s", err)
|
||||
}
|
||||
// ensure buffered events are sent before exit
|
||||
defer sentry.Flush(2 * time.Second)
|
||||
// Set up OpenTelemetry.
|
||||
otelShutdown, err := setupOTelSDK(ctx, cfg)
|
||||
if err != nil {
|
||||
panic(err.Error())
|
||||
}
|
||||
defer func() {
|
||||
err = errors.Join(err, otelShutdown(context.Background()))
|
||||
}()
|
||||
|
||||
var recordingHandler = recording.NewHandler(ctx, cfg)
|
||||
var transcodingHandler = transcoding.NewHandler(ctx, cfg)
|
||||
@@ -73,15 +71,11 @@ func main() {
|
||||
mux.Handle("/recording", http.StripPrefix("/recording", recordingHandler))
|
||||
mux.Handle("/transcoding", http.StripPrefix("/transcoding", transcodingHandler))
|
||||
|
||||
// wrap with Sentry HTTP handler if enabled
|
||||
// Wrap handler with OpenTelemetry instrumentation
|
||||
var handler http.Handler = mux
|
||||
if cfg.SentryDsn != "" {
|
||||
sentryHandler := sentryhttp.New(sentryhttp.Options{
|
||||
Repanic: true,
|
||||
WaitForDelivery: true,
|
||||
Timeout: 2 * time.Second,
|
||||
})
|
||||
handler = sentryHandler.Handle(mux)
|
||||
if cfg.OtelEnabled && cfg.OtelTracesEnabled {
|
||||
handler = otelhttp.NewHandler(handler, "http.server",
|
||||
otelhttp.WithServerName(cfg.OtelServiceName))
|
||||
}
|
||||
|
||||
server := &http.Server{
|
||||
|
||||
@@ -7,15 +7,27 @@ require (
|
||||
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/getsentry/sentry-go v0.31.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
|
||||
github.com/segmentio/kafka-go v0.4.48
|
||||
github.com/stretchr/testify v1.10.0
|
||||
github.com/stretchr/testify v1.11.1
|
||||
github.com/tus/tusd/v2 v2.6.0
|
||||
github.com/valyala/fasthttp v1.59.0
|
||||
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.63.0
|
||||
go.opentelemetry.io/otel v1.38.0
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.14.0
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp v1.38.0
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.14.0
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.38.0
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.38.0
|
||||
go.opentelemetry.io/otel/log v0.14.0
|
||||
go.opentelemetry.io/otel/sdk v1.38.0
|
||||
go.opentelemetry.io/otel/sdk/log v0.14.0
|
||||
go.opentelemetry.io/otel/sdk/metric v1.38.0
|
||||
go.opentelemetry.io/otel/trace v1.38.0
|
||||
go.uber.org/zap v1.27.0
|
||||
golang.org/x/exp v0.0.0-20250215185904-eff6e970281f
|
||||
gopkg.in/vansante/go-ffprobe.v2 v2.2.1
|
||||
@@ -38,14 +50,27 @@ require (
|
||||
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/aws/smithy-go v1.22.3 // indirect
|
||||
github.com/cenkalti/backoff/v5 v5.0.3 // indirect
|
||||
github.com/davecgh/go-spew v1.1.1 // indirect
|
||||
github.com/felixge/httpsnoop v1.0.4 // indirect
|
||||
github.com/go-logr/logr v1.4.3 // indirect
|
||||
github.com/go-logr/stdr v1.2.2 // indirect
|
||||
github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.2 // indirect
|
||||
github.com/klauspost/compress v1.18.0 // indirect
|
||||
github.com/pierrec/lz4/v4 v4.1.22 // indirect
|
||||
github.com/pmezard/go-difflib v1.0.0 // indirect
|
||||
github.com/valyala/bytebufferpool v1.0.0 // indirect
|
||||
go.opentelemetry.io/auto/sdk v1.1.0 // indirect
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0 // indirect
|
||||
go.opentelemetry.io/otel/metric v1.38.0 // indirect
|
||||
go.opentelemetry.io/proto/otlp v1.7.1 // indirect
|
||||
go.uber.org/multierr v1.11.0 // indirect
|
||||
golang.org/x/net v0.40.0 // indirect
|
||||
golang.org/x/sys v0.33.0 // indirect
|
||||
golang.org/x/text v0.25.0 // indirect
|
||||
golang.org/x/net v0.43.0 // indirect
|
||||
golang.org/x/sys v0.35.0 // indirect
|
||||
golang.org/x/text v0.28.0 // indirect
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20250825161204-c5933d9347a5 // indirect
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20250825161204-c5933d9347a5 // indirect
|
||||
google.golang.org/grpc v1.75.0 // indirect
|
||||
google.golang.org/protobuf v1.36.8 // indirect
|
||||
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||
)
|
||||
|
||||
@@ -38,49 +38,56 @@ github.com/aws/aws-sdk-go-v2/service/sts v1.33.14 h1:TzeR06UCMUq+KA3bDkujxK1GVGy
|
||||
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/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1xcsSM=
|
||||
github.com/cenkalti/backoff/v5 v5.0.3/go.mod h1:rkhZdG3JZukswDf7f0cwqPNk4K0sa+F97BxZthm/crw=
|
||||
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
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/getsentry/sentry-go v0.31.1 h1:ELVc0h7gwyhnXHDouXkhqTFSO5oslsRDk0++eyE0KJ4=
|
||||
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/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg=
|
||||
github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U=
|
||||
github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A=
|
||||
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
|
||||
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
|
||||
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
|
||||
github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE=
|
||||
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=
|
||||
github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY=
|
||||
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
|
||||
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
|
||||
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
|
||||
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
|
||||
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
|
||||
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
|
||||
github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.2 h1:8Tjv8EJ+pM1xP8mK6egEbD1OgnVTyacbefKhmbLhIhU=
|
||||
github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.2/go.mod h1:pkJQ2tZHJ0aFOVEEot6oZmaVEZcRme73eIFmhiVuRWs=
|
||||
github.com/kelseyhightower/envconfig v1.4.0 h1:Im6hONhd3pLkfDFsbRgu68RDNkGF1r3dvMUtDTo2cv8=
|
||||
github.com/kelseyhightower/envconfig v1.4.0/go.mod h1:cccZRl6mQpaq41TPp5QxidR+Sa3axMbJDNb//FQX6Gg=
|
||||
github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU=
|
||||
github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo=
|
||||
github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ=
|
||||
github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0=
|
||||
github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk=
|
||||
github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
|
||||
github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
|
||||
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
|
||||
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
|
||||
github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4=
|
||||
github.com/pierrec/lz4/v4 v4.1.22 h1:cKFw6uJDK+/gfw5BcDL0JL5aBsAFdsIT18eRtLj7VIU=
|
||||
github.com/pierrec/lz4/v4 v4.1.22/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4=
|
||||
github.com/pingcap/errors v0.11.4 h1:lFuQV/oaUMGcD2tqt+01ROSmJs75VG1ToEOkZIZ4nE4=
|
||||
github.com/pingcap/errors v0.11.4/go.mod h1:Oi8TUi2kEtXXLMJk9l1cGmz20kV3TaQ0usTwv5KuLY8=
|
||||
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/rogpeppe/go-internal v1.8.0 h1:FCbCCtXNOY3UtUuHUYaghJg4y7Fd14rXifAYUAtL9R8=
|
||||
github.com/rogpeppe/go-internal v1.8.0/go.mod h1:WmiCO8CzOY8rg0OYDC4/i/2WRWAB6poM+XZ2dLUbcbE=
|
||||
github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII=
|
||||
github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o=
|
||||
github.com/segmentio/kafka-go v0.4.48 h1:9jyu9CWK4W5W+SroCe8EffbrRZVqAOkuaLd/ApID4Vs=
|
||||
github.com/segmentio/kafka-go v0.4.48/go.mod h1:HjF6XbOKh0Pjlkr5GVZxt6CsjjwnmhVOfURM5KMd8qg=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw=
|
||||
github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
|
||||
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/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
|
||||
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
|
||||
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=
|
||||
@@ -96,6 +103,42 @@ github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gi
|
||||
github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU=
|
||||
github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E=
|
||||
github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY=
|
||||
go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA=
|
||||
go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A=
|
||||
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.63.0 h1:RbKq8BG0FI8OiXhBfcRtqqHcZcka+gU3cskNuf05R18=
|
||||
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.63.0/go.mod h1:h06DGIukJOevXaj/xrNjhi/2098RZzcLTbc0jDAUbsg=
|
||||
go.opentelemetry.io/otel v1.38.0 h1:RkfdswUDRimDg0m2Az18RKOsnI8UDzppJAtj01/Ymk8=
|
||||
go.opentelemetry.io/otel v1.38.0/go.mod h1:zcmtmQ1+YmQM9wrNsTGV/q/uyusom3P8RxwExxkZhjM=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.14.0 h1:QQqYw3lkrzwVsoEX0w//EhH/TCnpRdEenKBOOEIMjWc=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.14.0/go.mod h1:gSVQcr17jk2ig4jqJ2DX30IdWH251JcNAecvrqTxH1s=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp v1.38.0 h1:Oe2z/BCg5q7k4iXC3cqJxKYg0ieRiOqF0cecFYdPTwk=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp v1.38.0/go.mod h1:ZQM5lAJpOsKnYagGg/zV2krVqTtaVdYdDkhMoX6Oalg=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0 h1:GqRJVj7UmLjCVyVJ3ZFLdPRmhDUp2zFmQe3RHIOsw24=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0/go.mod h1:ri3aaHSmCTVYu2AWv44YMauwAQc0aqI9gHKIcSbI1pU=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0 h1:aTL7F04bJHUlztTsNGJ2l+6he8c+y/b//eR0jjjemT4=
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0/go.mod h1:kldtb7jDTeol0l3ewcmd8SDvx3EmIE7lyvqbasU3QC4=
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.14.0 h1:B/g+qde6Mkzxbry5ZZag0l7QrQBCtVm7lVjaLgmpje8=
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.14.0/go.mod h1:mOJK8eMmgW6ocDJn6Bn11CcZ05gi3P8GylBXEkZtbgA=
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.38.0 h1:wm/Q0GAAykXv83wzcKzGGqAnnfLFyFe7RslekZuv+VI=
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.38.0/go.mod h1:ra3Pa40+oKjvYh+ZD3EdxFZZB0xdMfuileHAm4nNN7w=
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.38.0 h1:kJxSDN4SgWWTjG/hPp3O7LCGLcHXFlvS2/FFOrwL+SE=
|
||||
go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.38.0/go.mod h1:mgIOzS7iZeKJdeB8/NYHrJ48fdGc71Llo5bJ1J4DWUE=
|
||||
go.opentelemetry.io/otel/log v0.14.0 h1:2rzJ+pOAZ8qmZ3DDHg73NEKzSZkhkGIua9gXtxNGgrM=
|
||||
go.opentelemetry.io/otel/log v0.14.0/go.mod h1:5jRG92fEAgx0SU/vFPxmJvhIuDU9E1SUnEQrMlJpOno=
|
||||
go.opentelemetry.io/otel/metric v1.38.0 h1:Kl6lzIYGAh5M159u9NgiRkmoMKjvbsKtYRwgfrA6WpA=
|
||||
go.opentelemetry.io/otel/metric v1.38.0/go.mod h1:kB5n/QoRM8YwmUahxvI3bO34eVtQf2i4utNVLr9gEmI=
|
||||
go.opentelemetry.io/otel/sdk v1.38.0 h1:l48sr5YbNf2hpCUj/FoGhW9yDkl+Ma+LrVl8qaM5b+E=
|
||||
go.opentelemetry.io/otel/sdk v1.38.0/go.mod h1:ghmNdGlVemJI3+ZB5iDEuk4bWA3GkTpW+DOoZMYBVVg=
|
||||
go.opentelemetry.io/otel/sdk/log v0.14.0 h1:JU/U3O7N6fsAXj0+CXz21Czg532dW2V4gG1HE/e8Zrg=
|
||||
go.opentelemetry.io/otel/sdk/log v0.14.0/go.mod h1:imQvII+0ZylXfKU7/wtOND8Hn4OpT3YUoIgqJVksUkM=
|
||||
go.opentelemetry.io/otel/sdk/log/logtest v0.14.0 h1:Ijbtz+JKXl8T2MngiwqBlPaHqc4YCaP/i13Qrow6gAM=
|
||||
go.opentelemetry.io/otel/sdk/log/logtest v0.14.0/go.mod h1:dCU8aEL6q+L9cYTqcVOk8rM9Tp8WdnHOPLiBgp0SGOA=
|
||||
go.opentelemetry.io/otel/sdk/metric v1.38.0 h1:aSH66iL0aZqo//xXzQLYozmWrXxyFkBJ6qT5wthqPoM=
|
||||
go.opentelemetry.io/otel/sdk/metric v1.38.0/go.mod h1:dg9PBnW9XdQ1Hd6ZnRz689CbtrUp0wMMs9iPcgT9EZA=
|
||||
go.opentelemetry.io/otel/trace v1.38.0 h1:Fxk5bKrDZJUH+AMyyIXGcFAPah0oRcT+LuNtJrmcNLE=
|
||||
go.opentelemetry.io/otel/trace v1.38.0/go.mod h1:j1P9ivuFsTceSWe1oY+EeW3sc+Pp42sO++GHkg4wwhs=
|
||||
go.opentelemetry.io/proto/otlp v1.7.1 h1:gTOMpGDb0WTBOP8JaO72iL3auEZhVmAQg4ipjOVAtj4=
|
||||
go.opentelemetry.io/proto/otlp v1.7.1/go.mod h1:b2rVh6rfI/s2pHWNlB7ILJcRALpcNDzKhACevjI+ZnE=
|
||||
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
|
||||
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
|
||||
go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0=
|
||||
@@ -115,8 +158,8 @@ golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug
|
||||
golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
|
||||
golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg=
|
||||
golang.org/x/net v0.17.0/go.mod h1:NxSsAGuq816PNPmqtQdLE42eU2Fs7NoRIZrHJAlaCOE=
|
||||
golang.org/x/net v0.40.0 h1:79Xs7wF06Gbdcg4kdCCIQArK11Z1hr5POQ6+fIYHNuY=
|
||||
golang.org/x/net v0.40.0/go.mod h1:y0hY0exeL2Pku80/zKK7tpntoX23cqL3Oa6njdgRtds=
|
||||
golang.org/x/net v0.43.0 h1:lat02VYK2j4aLzMzecihNvTlJNQUq316m2Mr9rnM6YE=
|
||||
golang.org/x/net v0.43.0/go.mod h1:vhO1fvI4dGsIjh73sWfUVjj3N7CA9WkKJNQm2svM6Jg=
|
||||
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
@@ -128,8 +171,8 @@ golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBc
|
||||
golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.33.0 h1:q3i8TbbEz+JRD9ywIRlyRAQbM0qF7hu24q3teo2hbuw=
|
||||
golang.org/x/sys v0.33.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
|
||||
golang.org/x/sys v0.35.0 h1:vz1N37gP5bs89s7He8XuIYXpyY0+QlsKmzipCbUtyxI=
|
||||
golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
|
||||
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=
|
||||
golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k=
|
||||
@@ -142,13 +185,23 @@ golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ=
|
||||
golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8=
|
||||
golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8=
|
||||
golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE=
|
||||
golang.org/x/text v0.25.0 h1:qVyWApTSYLk/drJRO5mDlNYskwQznZmkpV2c8q9zls4=
|
||||
golang.org/x/text v0.25.0/go.mod h1:WEdwpYrmk1qmdHvhkSTNPm3app7v4rsT8F2UD6+VHIA=
|
||||
golang.org/x/text v0.28.0 h1:rhazDwis8INMIwQ4tpjLDzUhx6RlXqZNPEM0huQojng=
|
||||
golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU=
|
||||
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
|
||||
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc=
|
||||
golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU=
|
||||
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk=
|
||||
gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E=
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20250825161204-c5933d9347a5 h1:BIRfGDEjiHRrk0QKZe3Xv2ieMhtgRGeLcZQ0mIVn4EY=
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20250825161204-c5933d9347a5/go.mod h1:j3QtIyytwqGr1JUDtYXwtMXWPKsEa5LtzIFN1Wn5WvE=
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20250825161204-c5933d9347a5 h1:eaY8u2EuxbRv7c3NiGK0/NedzVsCcV6hDuU5qPX5EGE=
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20250825161204-c5933d9347a5/go.mod h1:M4/wBTSeyLxupu3W3tJtOgB14jILAS/XWPSSa3TAlJc=
|
||||
google.golang.org/grpc v1.75.0 h1:+TW+dqTd2Biwe6KKfhE5JpiYIBWq865PhKGSXiivqt4=
|
||||
google.golang.org/grpc v1.75.0/go.mod h1:JtPAzKiq4v1xcAB2hydNlWI2RnF85XXcV0mhKXr2ecQ=
|
||||
google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc=
|
||||
google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
|
||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
|
||||
|
||||
@@ -25,7 +25,6 @@ import (
|
||||
|
||||
// Config represents configuration for the huly-stream application.
|
||||
type Config struct {
|
||||
SentryDsn string `split_words:"true" default:"" desc:"sentry dsn value"`
|
||||
LogLevel string `split_words:"true" default:"debug" desc:"sets log level for the application"`
|
||||
ServerSecret string `split_words:"true" default:"" desc:"server secret required to generate and verify tokens"`
|
||||
PprofEnabled bool `split_words:"true" default:"true" desc:"starts profile server on localhost:6060 if true"`
|
||||
@@ -40,6 +39,14 @@ type Config struct {
|
||||
|
||||
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"`
|
||||
|
||||
// OpenTelemetry configuration
|
||||
OtelEnabled bool `split_words:"true" default:"true" desc:"enable OpenTelemetry"`
|
||||
OtelServiceName string `split_words:"true" default:"stream" desc:"service name for OpenTelemetry"`
|
||||
OtelServiceVersion string `split_words:"true" default:"1.0.0" desc:"service version for OpenTelemetry"`
|
||||
OtelTracesEnabled bool `split_words:"true" default:"true" desc:"enable OpenTelemetry traces"`
|
||||
OtelMetricsEnabled bool `split_words:"true" default:"true" desc:"enable OpenTelemetry metrics"`
|
||||
OtelLogsEnabled bool `split_words:"true" default:"false" desc:"enable OpenTelemetry logs"`
|
||||
}
|
||||
|
||||
// FromEnv creates new Config from env
|
||||
|
||||
@@ -13,7 +13,7 @@
|
||||
// limitations under the License.
|
||||
//
|
||||
|
||||
package mediaconvert
|
||||
package executor
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
@@ -24,30 +24,18 @@ import (
|
||||
"sync"
|
||||
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/tracing"
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// CommandExecutor executes multiple commands in parallel
|
||||
type CommandExecutor interface {
|
||||
Execute(commands []*exec.Cmd) error
|
||||
}
|
||||
var tracer = otel.Tracer("executor")
|
||||
|
||||
type commandExecutor struct {
|
||||
logger *zap.Logger
|
||||
}
|
||||
|
||||
var _ CommandExecutor = (*commandExecutor)(nil)
|
||||
|
||||
// NewCommandExecutor creates a new instance of command executor
|
||||
func NewCommandExecutor(ctx context.Context) CommandExecutor {
|
||||
return &commandExecutor{
|
||||
logger: log.FromContext(ctx),
|
||||
}
|
||||
}
|
||||
|
||||
// Execute executes multiple commands in parallel
|
||||
func (e *commandExecutor) Execute(commands []*exec.Cmd) error {
|
||||
logger := e.logger
|
||||
// ExecuteCommands executes multiple commands in parallel
|
||||
func ExecuteCommands(ctx context.Context, commands []*exec.Cmd) error {
|
||||
logger := log.FromContext(ctx)
|
||||
errCh := make(chan error, len(commands))
|
||||
|
||||
var mu sync.Mutex
|
||||
@@ -58,6 +46,11 @@ func (e *commandExecutor) Execute(commands []*exec.Cmd) error {
|
||||
go func(cmd *exec.Cmd) {
|
||||
defer wg.Done()
|
||||
|
||||
_, span := tracer.Start(ctx, "cmd", trace.WithAttributes(
|
||||
attribute.String("command", cmd.String()),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var stdoutBuf = &bytes.Buffer{}
|
||||
var stderrBuf = &bytes.Buffer{}
|
||||
|
||||
@@ -71,6 +64,7 @@ func (e *commandExecutor) Execute(commands []*exec.Cmd) error {
|
||||
|
||||
logger.Info("run command", zap.String("cmd", cmd.String()))
|
||||
if err := cmd.Run(); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
errCh <- err
|
||||
|
||||
// Lock so only on goroutine can write to stdout/stderr at the same time
|
||||
@@ -13,7 +13,7 @@
|
||||
// limitations under the License.
|
||||
//
|
||||
|
||||
package mediaconvert_test
|
||||
package executor_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
@@ -21,8 +21,8 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/hcengineering/stream/internal/pkg/executor"
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/mediaconvert"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
@@ -55,9 +55,8 @@ func TestCommandExecutor_Execute_Success(t *testing.T) {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
ctx = log.WithFields(ctx)
|
||||
executor := mediaconvert.NewCommandExecutor(ctx)
|
||||
|
||||
err := executor.Execute(tt.commands)
|
||||
err := executor.ExecuteCommands(ctx, tt.commands)
|
||||
assert.NoError(t, err)
|
||||
})
|
||||
}
|
||||
@@ -107,9 +106,8 @@ func TestCommandExecutor_Execute_Error(t *testing.T) {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
ctx = log.WithFields(ctx)
|
||||
executor := mediaconvert.NewCommandExecutor(ctx)
|
||||
|
||||
err := executor.Execute(tt.commands)
|
||||
err := executor.ExecuteCommands(ctx, tt.commands)
|
||||
if tt.expectedError {
|
||||
assert.Error(t, err)
|
||||
} else {
|
||||
@@ -130,10 +128,9 @@ func TestCommandExecutor_Execute_Parallel(t *testing.T) {
|
||||
|
||||
ctx := context.Background()
|
||||
ctx = log.WithFields(ctx)
|
||||
executor := mediaconvert.NewCommandExecutor(ctx)
|
||||
|
||||
start := time.Now()
|
||||
err := executor.Execute(commands)
|
||||
err := executor.ExecuteCommands(ctx, commands)
|
||||
duration := time.Since(start)
|
||||
|
||||
assert.NoError(t, err)
|
||||
@@ -26,15 +26,21 @@ import (
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/hcengineering/stream/internal/pkg/config"
|
||||
"github.com/hcengineering/stream/internal/pkg/executor"
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/manifest"
|
||||
"github.com/hcengineering/stream/internal/pkg/storage"
|
||||
"github.com/hcengineering/stream/internal/pkg/token"
|
||||
"github.com/hcengineering/stream/internal/pkg/tracing"
|
||||
"github.com/hcengineering/stream/internal/pkg/uploader"
|
||||
"github.com/pkg/errors"
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.uber.org/zap"
|
||||
"gopkg.in/vansante/go-ffprobe.v2"
|
||||
)
|
||||
|
||||
var tracer = otel.Tracer("mediaconvert")
|
||||
|
||||
// Task represents transcoding task
|
||||
type Task struct {
|
||||
ID string
|
||||
@@ -97,14 +103,19 @@ func (p *Scheduler) start() {
|
||||
for range p.cfg.MaxParallelTranscodingCount {
|
||||
go func() {
|
||||
for task := range p.taskCh {
|
||||
p.processTask(p.ctx, task)
|
||||
err := tracing.WithSpan(p.ctx, tracer, "transcode", func(ctx context.Context) error {
|
||||
return p.processTask(ctx, task)
|
||||
})
|
||||
if err != nil {
|
||||
p.logger.Error("failed to process task", zap.Error(err))
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: add a factory pattern to process tasks by different media type
|
||||
func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
func (p *Scheduler) processTask(ctx context.Context, task *Task) error {
|
||||
var logger = p.logger.With(zap.String("task-id", task.ID))
|
||||
|
||||
logger.Debug("start")
|
||||
@@ -114,7 +125,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
var tokenString, err = token.NewToken(p.cfg.ServerSecret, task.Workspace, "stream")
|
||||
if err != nil {
|
||||
logger.Error("can not create token", zap.Error(err))
|
||||
return
|
||||
return errors.Wrapf(err, "can not create token")
|
||||
}
|
||||
|
||||
logger.Debug("phase 2: preparing fs")
|
||||
@@ -124,7 +135,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
err = os.MkdirAll(destinationFolder, os.ModePerm)
|
||||
if err != nil {
|
||||
logger.Error("can not create temporary folder", zap.Error(err))
|
||||
return
|
||||
return errors.Wrapf(err, "can not create temporary folder")
|
||||
}
|
||||
|
||||
defer func() {
|
||||
@@ -140,35 +151,37 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
if err != nil {
|
||||
logger.Error("can not create storage by url", zap.Error(err), zap.String("url", p.cfg.EndpointURL.String()))
|
||||
_ = os.RemoveAll(destinationFolder)
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
stat, err := remoteStorage.StatFile(ctx, task.Source)
|
||||
if err != nil {
|
||||
logger.Error("can not stat a file", zap.Error(err), zap.String("filepath", task.Source))
|
||||
_ = os.RemoveAll(destinationFolder)
|
||||
return
|
||||
return errors.Wrapf(err, "can not stat a file")
|
||||
}
|
||||
|
||||
if !IsSupportedMediaType(stat.Type) {
|
||||
logger.Info("unsupported media type", zap.String("type", stat.Type))
|
||||
_ = os.RemoveAll(destinationFolder)
|
||||
return
|
||||
return fmt.Errorf("unsupported media type: %s", stat.Type)
|
||||
}
|
||||
|
||||
if err = remoteStorage.GetFile(ctx, task.Source, sourceFilePath); err != nil {
|
||||
logger.Error("can not download a file", zap.Error(err), zap.String("filepath", task.Source))
|
||||
_ = os.RemoveAll(destinationFolder)
|
||||
// TODO: reschedule
|
||||
return
|
||||
return errors.Wrapf(err, "can not download a file")
|
||||
}
|
||||
|
||||
logger.Debug("phase 4: prepare to transcode")
|
||||
probe, err := ffprobe.ProbeURL(ctx, sourceFilePath)
|
||||
probe, err := tracing.WithSpanResult(ctx, tracer, "ffprobe", func(ctx context.Context) (*ffprobe.ProbeData, error) {
|
||||
return ffprobe.ProbeURL(ctx, sourceFilePath)
|
||||
})
|
||||
if err != nil {
|
||||
logger.Error("can not get probe for a file", zap.Error(err), zap.String("filepath", sourceFilePath))
|
||||
_ = os.RemoveAll(destinationFolder)
|
||||
return
|
||||
return errors.Wrapf(err, "can not get ffprobe")
|
||||
}
|
||||
|
||||
audioStream := probe.FirstAudioStream()
|
||||
@@ -180,7 +193,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
if videoStream == nil {
|
||||
logger.Error("no video stream found in the file", zap.String("filepath", sourceFilePath))
|
||||
_ = os.RemoveAll(destinationFolder)
|
||||
return
|
||||
return fmt.Errorf("no video stream found")
|
||||
}
|
||||
|
||||
logger.Debug("video stream found", zap.String("codec", videoStream.CodecName), zap.Int("width", videoStream.Width), zap.Int("height", videoStream.Height))
|
||||
@@ -221,7 +234,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
if err != nil {
|
||||
logger.Error("can not generate hls playlist", zap.String("out", p.cfg.OutputDir), zap.String("uploadID", opts.UploadID))
|
||||
_ = os.RemoveAll(destinationFolder)
|
||||
return
|
||||
return errors.Wrapf(err, "can not generate hls playlist")
|
||||
}
|
||||
|
||||
go uploader.Start()
|
||||
@@ -239,23 +252,19 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
if cmdErr != nil {
|
||||
logger.Error("can not create a new command", zap.Error(cmdErr), zap.Strings("args", args))
|
||||
go uploader.Cancel()
|
||||
return
|
||||
return errors.Wrapf(cmdErr, "can not create a new command")
|
||||
}
|
||||
cmds = append(cmds, cmd)
|
||||
if err = cmd.Start(); err != nil {
|
||||
logger.Error("can not start a command", zap.Error(err), zap.Strings("args", args))
|
||||
go uploader.Cancel()
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
logger.Debug("phase 7: wait for the result")
|
||||
for _, cmd := range cmds {
|
||||
if err = cmd.Wait(); err != nil {
|
||||
logger.Error("can not wait for command end ", zap.Error(err))
|
||||
go uploader.Cancel()
|
||||
return
|
||||
}
|
||||
execErr := tracing.WithSpan(ctx, tracer, "ffmpeg", func(ctx context.Context) error {
|
||||
return executor.ExecuteCommands(ctx, cmds)
|
||||
})
|
||||
if execErr != nil {
|
||||
logger.Error("can not wait for command end ", zap.Error(execErr))
|
||||
go uploader.Cancel()
|
||||
return errors.Wrapf(execErr, "can not execute command")
|
||||
}
|
||||
|
||||
logger.Debug("phase 8: schedule cleanup")
|
||||
@@ -293,6 +302,7 @@ func (p *Scheduler) processTask(ctx context.Context, task *Task) {
|
||||
logger.Error("can not patch the source file", zap.Error(err))
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// IsHLSSupportedVideoCodec checks whether the codec is supported by HLS
|
||||
|
||||
@@ -20,9 +20,12 @@ import (
|
||||
"sync"
|
||||
|
||||
"github.com/pkg/errors"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
|
||||
"github.com/hcengineering/stream/internal/pkg/sharedpipe"
|
||||
"github.com/hcengineering/stream/internal/pkg/storage"
|
||||
"github.com/hcengineering/stream/internal/pkg/tracing"
|
||||
"github.com/tus/tusd/v2/pkg/handler"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
@@ -46,6 +49,12 @@ var _ handler.LengthDeclarableUpload = (*Stream)(nil)
|
||||
|
||||
// WriteChunk is called when client sends a chunk of raw data
|
||||
func (w *Stream) WriteChunk(ctx context.Context, _ int64, src io.Reader) (int64, error) {
|
||||
ctx, span := tracer.Start(ctx, "WriteChunk", trace.WithAttributes(
|
||||
attribute.String("workspace", w.info.MetaData["workspace"]),
|
||||
attribute.String("upload_id", w.info.ID),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
w.logger.Debug("Write Chunk start", zap.Int64("offset", w.info.Offset))
|
||||
data, err := io.ReadAll(src)
|
||||
if err != nil {
|
||||
@@ -93,10 +102,17 @@ func (w *Stream) GetReader(ctx context.Context) (io.ReadCloser, error) {
|
||||
|
||||
// Terminate calls when upload has failed
|
||||
func (w *Stream) Terminate(ctx context.Context) error {
|
||||
ctx, span := tracer.Start(ctx, "Terminate", trace.WithAttributes(
|
||||
attribute.String("workspace", w.info.MetaData["workspace"]),
|
||||
attribute.String("upload_id", w.info.ID),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
w.logger.Debug("terminate upload")
|
||||
|
||||
// Close the writer first to signal EOF to all readers
|
||||
if err := w.writer.Close(); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
w.logger.Error("failed to close writer", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
@@ -133,10 +149,17 @@ func (w *Stream) ConcatUploads(ctx context.Context, partialUploads []handler.Upl
|
||||
|
||||
// FinishUpload calls when upload finished without errors on the client side
|
||||
func (w *Stream) FinishUpload(ctx context.Context) error {
|
||||
ctx, span := tracer.Start(ctx, "FinishUpload", trace.WithAttributes(
|
||||
attribute.String("workspace", w.info.MetaData["workspace"]),
|
||||
attribute.String("upload_id", w.info.ID),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
w.logger.Debug("finish upload")
|
||||
|
||||
// Close the writer first to signal EOF to all readers
|
||||
if err := w.writer.Close(); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
w.logger.Error("failed to close writer", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
@@ -151,6 +174,7 @@ 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))
|
||||
tracing.RecordError(span, err)
|
||||
completeErr = err
|
||||
return
|
||||
}
|
||||
|
||||
@@ -24,10 +24,12 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/hcengineering/stream/internal/pkg/config"
|
||||
"github.com/hcengineering/stream/internal/pkg/executor"
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/manifest"
|
||||
"github.com/hcengineering/stream/internal/pkg/storage"
|
||||
"github.com/hcengineering/stream/internal/pkg/token"
|
||||
"github.com/hcengineering/stream/internal/pkg/tracing"
|
||||
"github.com/hcengineering/stream/internal/pkg/uploader"
|
||||
"github.com/pkg/errors"
|
||||
"go.uber.org/zap"
|
||||
@@ -115,7 +117,9 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er
|
||||
}
|
||||
|
||||
logger.Debug("phase 4: prepare to transcode")
|
||||
probe, err := ffprobe.ProbeURL(ctx, sourceFilePath)
|
||||
probe, err := tracing.WithSpanResult(ctx, tracer, "ffprobe", func(spanCtx context.Context) (*ffprobe.ProbeData, error) {
|
||||
return ffprobe.ProbeURL(spanCtx, sourceFilePath)
|
||||
})
|
||||
if err != nil {
|
||||
logger.Error("can not get ffprobe", zap.Error(err), zap.String("filepath", sourceFilePath))
|
||||
return nil, errors.Wrapf(err, "can not get ffprobe")
|
||||
@@ -196,8 +200,10 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, er
|
||||
cmds = append(cmds, cmd)
|
||||
}
|
||||
|
||||
executor := NewCommandExecutor(ctx)
|
||||
if execErr := executor.Execute(cmds); execErr != nil {
|
||||
execErr := tracing.WithSpan(ctx, tracer, "ffmpeg", func(ctx context.Context) error {
|
||||
return executor.ExecuteCommands(ctx, cmds)
|
||||
})
|
||||
if execErr != nil {
|
||||
uploader.Cancel()
|
||||
return nil, errors.Wrapf(execErr, "can not execute command")
|
||||
}
|
||||
|
||||
@@ -21,11 +21,15 @@ import (
|
||||
"github.com/hcengineering/stream/internal/pkg/config"
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/mediaconvert"
|
||||
"github.com/hcengineering/stream/internal/pkg/tracing"
|
||||
"github.com/pkg/errors"
|
||||
"github.com/segmentio/kafka-go"
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
var tracer = otel.Tracer("worker")
|
||||
|
||||
// Worker is a queue processing worker
|
||||
type Worker struct {
|
||||
logger *zap.Logger
|
||||
@@ -109,7 +113,9 @@ func (w *Worker) processMessage(ctx context.Context, msg kafka.Message, logger *
|
||||
}
|
||||
|
||||
transcoder := mediaconvert.NewTranscoder(ctx, w.cfg)
|
||||
res, err := transcoder.Transcode(ctx, &task)
|
||||
res, err := tracing.WithSpanResult(ctx, tracer, "transcode", func(spanCtx context.Context) (*mediaconvert.TaskResult, error) {
|
||||
return transcoder.Transcode(spanCtx, &task)
|
||||
})
|
||||
|
||||
if res != nil {
|
||||
result := TranscodeResult{
|
||||
|
||||
@@ -29,11 +29,17 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/tracing"
|
||||
"github.com/pkg/errors"
|
||||
"github.com/valyala/fasthttp"
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
var tracer = otel.Tracer("storage.datalake")
|
||||
|
||||
type uploadResult struct {
|
||||
key string
|
||||
error string
|
||||
@@ -84,9 +90,15 @@ func getObjectKeyFromPath(s string) string {
|
||||
}
|
||||
|
||||
// PutFile uploads file to the datalake
|
||||
func (d *DatalakeStorage) PutFile(ctx context.Context, fileName string, options PutOptions) error {
|
||||
func (d *DatalakeStorage) PutFile(ctx context.Context, filename string, options PutOptions) error {
|
||||
ctx, span := tracer.Start(ctx, "datalake.put_file", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", getObjectKeyFromPath(filename)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
// #nosec
|
||||
file, err := os.Open(fileName)
|
||||
file, err := os.Open(filename)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -94,8 +106,8 @@ func (d *DatalakeStorage) PutFile(ctx context.Context, fileName string, options
|
||||
_ = file.Close()
|
||||
}()
|
||||
|
||||
var objectKey = getObjectKeyFromPath(fileName)
|
||||
var logger = d.logger.With(zap.String("upload", d.workspace), zap.String("fileName", fileName))
|
||||
var objectKey = getObjectKeyFromPath(filename)
|
||||
var logger = d.logger.With(zap.String("upload", d.workspace), zap.String("fileName", filename))
|
||||
|
||||
logger.Debug("start")
|
||||
|
||||
@@ -104,16 +116,19 @@ func (d *DatalakeStorage) PutFile(ctx context.Context, fileName string, options
|
||||
|
||||
part, err := createFormFile(writer, "file", objectKey, getContentType(objectKey))
|
||||
if err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
return errors.Wrapf(err, "failed to create form file")
|
||||
}
|
||||
|
||||
_, err = io.Copy(part, file)
|
||||
if err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
return errors.Wrapf(err, "failed to copy file data")
|
||||
}
|
||||
|
||||
err = writer.Close()
|
||||
if err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
return errors.Wrapf(err, "failed to close multipart writer")
|
||||
}
|
||||
|
||||
@@ -133,32 +148,39 @@ func (d *DatalakeStorage) PutFile(ctx context.Context, fileName string, options
|
||||
req.SetBody(body.Bytes())
|
||||
|
||||
if err := d.client.Do(req, res); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "upload failed", res)
|
||||
return errors.Wrapf(err, "upload failed")
|
||||
}
|
||||
|
||||
var result []uploadResult
|
||||
if err := json.Unmarshal(res.Body(), &result); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
return errors.Wrapf(err, "parse error")
|
||||
}
|
||||
|
||||
for _, res := range result {
|
||||
if res.error != "" {
|
||||
tracing.RecordError(span, err)
|
||||
return fmt.Errorf("upload error: %v %v", res.key, res.error)
|
||||
}
|
||||
}
|
||||
|
||||
logger.Debug("uploaded")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteFile deletes file from the datalake
|
||||
func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error {
|
||||
var logger = d.logger.With(zap.String("delete", d.workspace), zap.String("fileName", fileName))
|
||||
func (d *DatalakeStorage) DeleteFile(ctx context.Context, filename string) error {
|
||||
ctx, span := tracer.Start(ctx, "datalake.delete_file", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", getObjectKeyFromPath(filename)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("delete", d.workspace), zap.String("fileName", filename))
|
||||
logger.Debug("start")
|
||||
|
||||
var objectKey = getObjectKeyFromPath(fileName)
|
||||
var objectKey = getObjectKeyFromPath(filename)
|
||||
|
||||
req := fasthttp.AcquireRequest()
|
||||
defer fasthttp.ReleaseRequest(req)
|
||||
@@ -171,22 +193,28 @@ func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error
|
||||
req.Header.Add("Authorization", "Bearer "+d.token)
|
||||
|
||||
if err := d.client.Do(req, res); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "delete failed", res)
|
||||
return errors.Wrapf(err, "delete failed")
|
||||
}
|
||||
|
||||
if err := okResponse(res); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", res)
|
||||
return err
|
||||
}
|
||||
|
||||
logger.Debug("deleted")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// PatchMeta patches metadata for the object
|
||||
func (d *DatalakeStorage) PatchMeta(ctx context.Context, filename string, md *Metadata) error {
|
||||
ctx, span := tracer.Start(ctx, "datalake.patch_meta", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", getObjectKeyFromPath(filename)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("patch meta", d.workspace), zap.String("fileName", filename))
|
||||
logger.Debug("start")
|
||||
defer logger.Debug("finished")
|
||||
@@ -211,11 +239,13 @@ func (d *DatalakeStorage) PatchMeta(ctx context.Context, filename string, md *Me
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
@@ -225,6 +255,12 @@ func (d *DatalakeStorage) PatchMeta(ctx context.Context, filename string, md *Me
|
||||
|
||||
// GetMeta gets metadata related to the object
|
||||
func (d *DatalakeStorage) GetMeta(ctx context.Context, filename string) (*Metadata, error) {
|
||||
ctx, span := tracer.Start(ctx, "datalake.put_file", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", getObjectKeyFromPath(filename)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("get meta", d.workspace), zap.String("fileName", filename))
|
||||
logger.Debug("start")
|
||||
|
||||
@@ -240,11 +276,13 @@ func (d *DatalakeStorage) GetMeta(ctx context.Context, filename string) (*Metada
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return nil, err
|
||||
}
|
||||
@@ -257,6 +295,12 @@ func (d *DatalakeStorage) GetMeta(ctx context.Context, filename string) (*Metada
|
||||
|
||||
// GetFile gets file from the storage
|
||||
func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination string) error {
|
||||
ctx, span := tracer.Start(ctx, "datalake.get_file", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", getObjectKeyFromPath(filename)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("get", d.workspace), zap.String("fileName", filename), zap.String("destination", destination))
|
||||
logger.Debug("start")
|
||||
|
||||
@@ -272,11 +316,13 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
@@ -284,6 +330,7 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str
|
||||
// #nosec
|
||||
file, err := os.Create(destination)
|
||||
if err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logger.Debug("can't create a file", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
@@ -291,12 +338,14 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str
|
||||
_ = file.Close()
|
||||
}()
|
||||
if err = resp.BodyWriteTo(file); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logger.Debug("can't write to file", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
|
||||
stat, err := os.Stat(destination)
|
||||
if err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logger.Error("can't stat the file", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
@@ -307,6 +356,12 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str
|
||||
|
||||
// StatFile gets file stat from the storage
|
||||
func (d *DatalakeStorage) StatFile(ctx context.Context, filename string) (*BlobInfo, error) {
|
||||
ctx, span := tracer.Start(ctx, "datalake.stat_file", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", getObjectKeyFromPath(filename)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("head", d.workspace), zap.String("fileName", filename))
|
||||
logger.Debug("start")
|
||||
|
||||
@@ -322,11 +377,13 @@ func (d *DatalakeStorage) StatFile(ctx context.Context, filename string) (*BlobI
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return nil, err
|
||||
}
|
||||
@@ -342,6 +399,12 @@ func (d *DatalakeStorage) StatFile(ctx context.Context, filename string) (*BlobI
|
||||
|
||||
// SetParent updates blob parent reference
|
||||
func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string) error {
|
||||
ctx, span := tracer.Start(ctx, "datalake.set_parent", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", getObjectKeyFromPath(filename)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("workspace", d.workspace), zap.String("fileName", filename), zap.String("parent", parent))
|
||||
|
||||
logger.Debug("start")
|
||||
@@ -366,6 +429,7 @@ func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string
|
||||
}
|
||||
|
||||
if err := json.NewEncoder(req.BodyWriter()).Encode(body); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logger.Debug("can not encode body", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
@@ -374,11 +438,13 @@ func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
@@ -388,6 +454,12 @@ func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string
|
||||
|
||||
// MultipartUploadStart creates a new multipart upload
|
||||
func (d *DatalakeStorage) MultipartUploadStart(ctx context.Context, objectName, contentType string) (string, error) {
|
||||
ctx, span := tracer.Start(ctx, "datalake.multipart_upload_start", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", objectName),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("workspace", d.workspace), zap.String("objectName", objectName))
|
||||
url := fmt.Sprintf("%v/upload/multipart/%v/%v", d.baseURL, d.workspace, objectName)
|
||||
|
||||
@@ -402,11 +474,13 @@ func (d *DatalakeStorage) MultipartUploadStart(ctx context.Context, objectName,
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return "", err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return "", err
|
||||
}
|
||||
@@ -421,6 +495,13 @@ func (d *DatalakeStorage) MultipartUploadStart(ctx context.Context, objectName,
|
||||
|
||||
// MultipartUploadPart uploads a part of a multipart upload
|
||||
func (d *DatalakeStorage) MultipartUploadPart(ctx context.Context, objectName, uploadID string, partNumber int, data []byte) (*MultipartPart, error) {
|
||||
ctx, span := tracer.Start(ctx, "datalake.multipart_upload_part", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", objectName),
|
||||
attribute.Int("size", len(data)),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("workspace", d.workspace), zap.String("uploadID", uploadID), zap.Int("partNumber", partNumber))
|
||||
params := url.Values{}
|
||||
params.Add("uploadId", uploadID)
|
||||
@@ -440,11 +521,13 @@ func (d *DatalakeStorage) MultipartUploadPart(ctx context.Context, objectName, u
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return nil, err
|
||||
}
|
||||
@@ -457,6 +540,12 @@ func (d *DatalakeStorage) MultipartUploadPart(ctx context.Context, objectName, u
|
||||
|
||||
// MultipartUploadComplete completes a multipart upload
|
||||
func (d *DatalakeStorage) MultipartUploadComplete(ctx context.Context, objectName, uploadID string, parts []MultipartPart) error {
|
||||
ctx, span := tracer.Start(ctx, "datalake.multipart_upload_complete", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", objectName),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("workspace", d.workspace), zap.String("uploadID", uploadID), zap.String("objectName", objectName))
|
||||
params := url.Values{}
|
||||
params.Add("uploadId", uploadID)
|
||||
@@ -483,11 +572,13 @@ func (d *DatalakeStorage) MultipartUploadComplete(ctx context.Context, objectNam
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
@@ -497,6 +588,12 @@ func (d *DatalakeStorage) MultipartUploadComplete(ctx context.Context, objectNam
|
||||
|
||||
// MultipartUploadCancel cancels a multipart upload
|
||||
func (d *DatalakeStorage) MultipartUploadCancel(ctx context.Context, objectName, uploadID string) error {
|
||||
ctx, span := tracer.Start(ctx, "datalake.multipart_upload_cancel", trace.WithAttributes(
|
||||
attribute.String("workspace", d.workspace),
|
||||
attribute.String("object_key", objectName),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
var logger = d.logger.With(zap.String("workspace", d.workspace), zap.String("uploadID", uploadID))
|
||||
params := url.Values{}
|
||||
params.Add("uploadId", uploadID)
|
||||
@@ -512,11 +609,13 @@ func (d *DatalakeStorage) MultipartUploadCancel(ctx context.Context, objectName,
|
||||
defer fasthttp.ReleaseResponse(resp)
|
||||
|
||||
if err := d.client.Do(req, resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "request failed", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
if err := okResponse(resp); err != nil {
|
||||
tracing.RecordError(span, err)
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
// 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 tracing
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
)
|
||||
|
||||
// RecordError records an error into the provided span and sets its status to Error.
|
||||
func RecordError(span trace.Span, err error) {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
}
|
||||
|
||||
// WithSpan wraps a function execution within a span.
|
||||
func WithSpan(ctx context.Context, tracer trace.Tracer, name string, fn func(ctx context.Context) error) error {
|
||||
ctx, span := tracer.Start(ctx, name)
|
||||
defer span.End()
|
||||
|
||||
err := fn(ctx)
|
||||
if err != nil {
|
||||
RecordError(span, err)
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
// WithSpanResult wraps a function returning result execution within a span.
|
||||
func WithSpanResult[Result any](ctx context.Context, tracer trace.Tracer, name string, fn func(ctx context.Context) (Result, error)) (Result, error) {
|
||||
ctx, span := tracer.Start(ctx, name)
|
||||
defer span.End()
|
||||
|
||||
result, err := fn(ctx)
|
||||
if err != nil {
|
||||
RecordError(span, err)
|
||||
}
|
||||
|
||||
return result, err
|
||||
}
|
||||
@@ -28,10 +28,13 @@ import (
|
||||
"github.com/hcengineering/stream/internal/pkg/log"
|
||||
"github.com/hcengineering/stream/internal/pkg/storage"
|
||||
"github.com/pkg/errors"
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.uber.org/zap"
|
||||
"k8s.io/utils/inotify"
|
||||
)
|
||||
|
||||
var tracer = otel.Tracer("uploader")
|
||||
|
||||
// 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
|
||||
@@ -93,7 +96,7 @@ func New(ctx context.Context, s storage.Storage, opts Options) Uploader {
|
||||
|
||||
res.logger.Sugar().Debugf("uploader config is %v", opts)
|
||||
|
||||
res.uploadCtx, res.uploadCancel = context.WithCancel(context.Background())
|
||||
res.uploadCtx, res.uploadCancel = context.WithCancel(ctx)
|
||||
|
||||
err := os.MkdirAll(opts.Dir, os.ModePerm)
|
||||
if err != nil {
|
||||
@@ -105,15 +108,18 @@ func New(ctx context.Context, s storage.Storage, opts Options) Uploader {
|
||||
|
||||
func (u *uploaderImpl) Stop() {
|
||||
u.logger.Info("stopping upload")
|
||||
u.stop(false)
|
||||
u.stop(u.uploadCtx, false)
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) Cancel() {
|
||||
u.logger.Info("canceling upload")
|
||||
u.stop(true)
|
||||
u.stop(u.uploadCtx, true)
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) scanFiles() {
|
||||
func (u *uploaderImpl) scanFiles(ctx context.Context) {
|
||||
_, span := tracer.Start(ctx, "scanFiles")
|
||||
defer span.End()
|
||||
|
||||
logger := u.logger.With(zap.String("dir", u.options.Dir))
|
||||
|
||||
logger.Info("scan files")
|
||||
@@ -148,13 +154,16 @@ func (u *uploaderImpl) scanFiles() {
|
||||
logger.Info("scan complete", zap.Int("count", count))
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) stop(rollback bool) {
|
||||
func (u *uploaderImpl) stop(ctx context.Context, rollback bool) {
|
||||
ctx, span := tracer.Start(ctx, "stop")
|
||||
defer span.End()
|
||||
|
||||
// Stop watching for new files
|
||||
close(u.watcherStopCh)
|
||||
<-u.watcherDoneCh
|
||||
|
||||
// Scan remaining files in the directory
|
||||
u.scanFiles()
|
||||
u.scanFiles(ctx)
|
||||
|
||||
// Close filesCh so no new files added
|
||||
close(u.filesCh)
|
||||
@@ -165,7 +174,7 @@ func (u *uploaderImpl) stop(rollback bool) {
|
||||
|
||||
// Perform rollback
|
||||
if rollback {
|
||||
u.uploadRollback()
|
||||
u.uploadRollback(ctx)
|
||||
}
|
||||
|
||||
u.uploadCancel()
|
||||
@@ -188,7 +197,10 @@ func (u *uploaderImpl) stop(rollback bool) {
|
||||
u.logger.Debug("stopped", zap.Bool("rollback", rollback))
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) uploadRollback() {
|
||||
func (u *uploaderImpl) uploadRollback(ctx context.Context) {
|
||||
_, span := tracer.Start(ctx, "uploadRollback")
|
||||
defer span.End()
|
||||
|
||||
u.logger.Debug("starting rollback...")
|
||||
|
||||
// Create a separate worker pool for rollback
|
||||
@@ -226,7 +238,7 @@ func (u *uploaderImpl) Start() {
|
||||
|
||||
<-watcherReady
|
||||
|
||||
u.scanFiles()
|
||||
u.scanFiles(u.uploadCtx)
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) startWorkers() {
|
||||
|
||||
Reference in New Issue
Block a user