From bca70dea13caa9136128bcd05c48ac7013a831f7 Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Fri, 24 Oct 2025 00:06:47 +0700 Subject: [PATCH] Open telemetry support Signed-off-by: Alexander Onnikov --- cmd/stream/main.go | 32 ++--- go.mod | 35 ++++- go.sum | 93 +++++++++++--- internal/pkg/config/config.go | 9 +- .../{mediaconvert => executor}/executor.go | 36 +++--- .../executor_test.go | 13 +- internal/pkg/mediaconvert/scheduler.go | 58 +++++---- internal/pkg/mediaconvert/stream.go | 24 ++++ internal/pkg/mediaconvert/transcoder.go | 12 +- internal/pkg/queue/worker.go | 8 +- internal/pkg/storage/datalake.go | 121 ++++++++++++++++-- internal/pkg/tracing/tracing.go | 53 ++++++++ internal/pkg/uploader/uploader.go | 30 +++-- 13 files changed, 402 insertions(+), 122 deletions(-) rename internal/pkg/{mediaconvert => executor}/executor.go (77%) rename internal/pkg/{mediaconvert => executor}/executor_test.go (89%) create mode 100644 internal/pkg/tracing/tracing.go diff --git a/cmd/stream/main.go b/cmd/stream/main.go index 37174f32af..3b00fd7d0c 100644 --- a/cmd/stream/main.go +++ b/cmd/stream/main.go @@ -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{ diff --git a/go.mod b/go.mod index 701291bda0..4c0266ff8a 100644 --- a/go.mod +++ b/go.mod @@ -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 ) diff --git a/go.sum b/go.sum index 3ca0c0a4f2..7a61312aa8 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/internal/pkg/config/config.go b/internal/pkg/config/config.go index ca0e4aef22..61cc586a52 100644 --- a/internal/pkg/config/config.go +++ b/internal/pkg/config/config.go @@ -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 diff --git a/internal/pkg/mediaconvert/executor.go b/internal/pkg/executor/executor.go similarity index 77% rename from internal/pkg/mediaconvert/executor.go rename to internal/pkg/executor/executor.go index 6385160aa7..e0f2cc90d9 100644 --- a/internal/pkg/mediaconvert/executor.go +++ b/internal/pkg/executor/executor.go @@ -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 diff --git a/internal/pkg/mediaconvert/executor_test.go b/internal/pkg/executor/executor_test.go similarity index 89% rename from internal/pkg/mediaconvert/executor_test.go rename to internal/pkg/executor/executor_test.go index 3f30a82cf9..4e8bf3a4d3 100644 --- a/internal/pkg/mediaconvert/executor_test.go +++ b/internal/pkg/executor/executor_test.go @@ -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) diff --git a/internal/pkg/mediaconvert/scheduler.go b/internal/pkg/mediaconvert/scheduler.go index 8235f44b02..d41998f75d 100644 --- a/internal/pkg/mediaconvert/scheduler.go +++ b/internal/pkg/mediaconvert/scheduler.go @@ -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 diff --git a/internal/pkg/mediaconvert/stream.go b/internal/pkg/mediaconvert/stream.go index fbb671c911..d32e05c2a6 100644 --- a/internal/pkg/mediaconvert/stream.go +++ b/internal/pkg/mediaconvert/stream.go @@ -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 } diff --git a/internal/pkg/mediaconvert/transcoder.go b/internal/pkg/mediaconvert/transcoder.go index a3fc1b738d..c0b4f9bc0a 100644 --- a/internal/pkg/mediaconvert/transcoder.go +++ b/internal/pkg/mediaconvert/transcoder.go @@ -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") } diff --git a/internal/pkg/queue/worker.go b/internal/pkg/queue/worker.go index a80b4c702a..4fc0802f0b 100644 --- a/internal/pkg/queue/worker.go +++ b/internal/pkg/queue/worker.go @@ -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{ diff --git a/internal/pkg/storage/datalake.go b/internal/pkg/storage/datalake.go index 9207a88ed9..1101f49f51 100644 --- a/internal/pkg/storage/datalake.go +++ b/internal/pkg/storage/datalake.go @@ -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 } diff --git a/internal/pkg/tracing/tracing.go b/internal/pkg/tracing/tracing.go new file mode 100644 index 0000000000..e9afbc55dd --- /dev/null +++ b/internal/pkg/tracing/tracing.go @@ -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 +} diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go index 2bdee38c76..78df9f3d2b 100644 --- a/internal/pkg/uploader/uploader.go +++ b/internal/pkg/uploader/uploader.go @@ -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() {