From c8fb7f3cc3edb300613c651512075fa9cd12d3dc Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sun, 4 May 2025 14:57:19 -0700 Subject: [PATCH] media: add SignMP4 tracing, fix channel leak (#126) * tracing: add a lot more instrumentation around SignMP4 tracing * tracing: track signing duration for external signing too * segchanman: time out segment send * concat: okay yeah i see why this was a leak now * webrtc_playback: actually unsubscribe here too --- go.mod | 7 +------ go.sum | 10 ---------- pkg/media/concat.go | 4 ++-- pkg/media/media_signer.go | 22 ++++++++++++++++++++++ pkg/media/media_signer_ext.go | 12 +++++++++++- pkg/media/segchanman/segchanman.go | 13 +++++++++++-- pkg/media/segmenter.go | 6 +++++- pkg/media/webrtc_playback.go | 4 ++-- pkg/spmetrics/spmetrics.go | 6 ++++++ 9 files changed, 60 insertions(+), 24 deletions(-) diff --git a/go.mod b/go.mod index 1b9c28486..9a37c6faf 100644 --- a/go.mod +++ b/go.mod @@ -49,13 +49,7 @@ require ( gitlab.com/gitlab-org/release-cli v0.18.0 go.opentelemetry.io/otel v1.35.0 go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.35.0 - go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.11.0 - go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.35.0 - go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.35.0 - go.opentelemetry.io/otel/log v0.11.0 go.opentelemetry.io/otel/sdk v1.35.0 - go.opentelemetry.io/otel/sdk/log v0.11.0 - go.opentelemetry.io/otel/sdk/metric v1.35.0 go.uber.org/goleak v1.3.0 golang.org/x/exp v0.0.0-20240909161429-701f63a606c0 golang.org/x/image v0.22.0 @@ -244,6 +238,7 @@ require ( go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.49.0 // indirect go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.35.0 // indirect go.opentelemetry.io/otel/metric v1.35.0 // indirect + go.opentelemetry.io/otel/sdk/metric v1.35.0 // indirect go.opentelemetry.io/otel/trace v1.35.0 // indirect go.opentelemetry.io/proto/otlp v1.5.0 // indirect go.uber.org/atomic v1.11.0 // indirect diff --git a/go.sum b/go.sum index a2ff16d94..648f141b9 100644 --- a/go.sum +++ b/go.sum @@ -657,20 +657,10 @@ go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.35.0 h1:1fTNlAIJZGWLP5FVu0f go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.35.0/go.mod h1:zjPK58DtkqQFn+YUMbx0M2XV3QgKU0gS9LeGohREyK4= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.35.0 h1:m639+BofXTvcY1q8CGs4ItwQarYtJPOWmVobfM1HpVI= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.35.0/go.mod h1:LjReUci/F4BUyv+y4dwnq3h/26iNOeC3wAIqgvTIZVo= -go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.11.0 h1:k6KdfZk72tVW/QVZf60xlDziDvYAePj5QHwoQvrB2m8= -go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.11.0/go.mod h1:5Y3ZJLqzi/x/kYtrSrPSx7TFI/SGsL7q2kME027tH6I= -go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.35.0 h1:PB3Zrjs1sG1GBX51SXyTSoOTqcDglmsk7nT6tkKPb/k= -go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.35.0/go.mod h1:U2R3XyVPzn0WX7wOIypPuptulsMcPDPs/oiSVOMVnHY= -go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.35.0 h1:T0Ec2E+3YZf5bgTNQVet8iTDW7oIk03tXHq+wkwIDnE= -go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.35.0/go.mod h1:30v2gqH+vYGJsesLWFov8u47EpYTcIQcBjKpI6pJThg= -go.opentelemetry.io/otel/log v0.11.0 h1:c24Hrlk5WJ8JWcwbQxdBqxZdOK7PcP/LFtOtwpDTe3Y= -go.opentelemetry.io/otel/log v0.11.0/go.mod h1:U/sxQ83FPmT29trrifhQg+Zj2lo1/IPN1PF6RTFqdwc= go.opentelemetry.io/otel/metric v1.35.0 h1:0znxYu2SNyuMSQT4Y9WDWej0VpcsxkuklLa4/siN90M= go.opentelemetry.io/otel/metric v1.35.0/go.mod h1:nKVFgxBZ2fReX6IlyW28MgZojkoAkJGaE8CpgeAU3oE= go.opentelemetry.io/otel/sdk v1.35.0 h1:iPctf8iprVySXSKJffSS79eOjl9pvxV9ZqOWT0QejKY= go.opentelemetry.io/otel/sdk v1.35.0/go.mod h1:+ga1bZliga3DxJ3CQGg3updiaAJoNECOgJREo9KHGQg= -go.opentelemetry.io/otel/sdk/log v0.11.0 h1:7bAOpjpGglWhdEzP8z0VXc4jObOiDEwr3IYbhBnjk2c= -go.opentelemetry.io/otel/sdk/log v0.11.0/go.mod h1:dndLTxZbwBstZoqsJB3kGsRPkpAgaJrWfQg3lhlHFFY= go.opentelemetry.io/otel/sdk/metric v1.35.0 h1:1RriWBmCKgkeHEhM7a2uMjMUfP7MsOF5JpUCaEqEI9o= go.opentelemetry.io/otel/sdk/metric v1.35.0/go.mod h1:is6XYCUMpcKi+ZsOvfluY5YstFnhW0BidkR+gL+qN+w= go.opentelemetry.io/otel/trace v1.35.0 h1:dPpEfJu1sDIqruz7BHFG3c7528f6ddfSWfFDVt/xgMs= diff --git a/pkg/media/concat.go b/pkg/media/concat.go index cb45803a4..bec1b9572 100644 --- a/pkg/media/concat.go +++ b/pkg/media/concat.go @@ -126,12 +126,12 @@ func ConcatStream(ctx context.Context, pipeline *gst.Pipeline, user string, rend // them in a pipe so that we don't miss any in between iterations of the output allFiles := make(chan []byte, 1024) go func() { + ch := streamer.SubscribeSegment(ctx, user, rendition) + defer streamer.UnsubscribeSegment(ctx, user, rendition, ch) for { - ch := streamer.SubscribeSegment(ctx, user, rendition) select { case <-ctx.Done(): log.Debug(ctx, "exiting segment reader") - streamer.UnsubscribeSegment(ctx, user, rendition, ch) return case file := <-ch: log.Debug(ctx, "got segment", "file", file.Filepath) diff --git a/pkg/media/media_signer.go b/pkg/media/media_signer.go index 84e526197..240cfe37d 100644 --- a/pkg/media/media_signer.go +++ b/pkg/media/media_signer.go @@ -9,6 +9,7 @@ import ( "fmt" "io" "path/filepath" + "time" "git.stream.place/streamplace/c2pa-go/pkg/c2pa" "go.opentelemetry.io/otel" @@ -18,11 +19,13 @@ import ( "stream.place/streamplace/pkg/crypto/aqpub" "stream.place/streamplace/pkg/crypto/signers" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/spmetrics" ) type MediaSigner interface { SignMP4(ctx context.Context, input io.ReadSeeker, start int64) ([]byte, error) Pub() aqpub.Pub + Streamer() string } type MediaSignerLocal struct { @@ -80,7 +83,12 @@ func MakeMediaSigner(ctx context.Context, cli *config.CLI, streamer string, sign }, nil } +func (ms *MediaSignerLocal) Streamer() string { + return ms.StreamerName +} + func (ms *MediaSignerLocal) SignMP4(ctx context.Context, input io.ReadSeeker, start int64) ([]byte, error) { + startTime := time.Now() ctx, span := otel.Tracer("signer").Start(ctx, "SignMP4") defer span.End() title := "livestream" @@ -109,6 +117,7 @@ func (ms *MediaSignerLocal) SignMP4(ctx context.Context, input io.ReadSeeker, st }, }, } + ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_MarshalManifest") manifestBs, err := json.Marshal(mani) if err != nil { return nil, fmt.Errorf("failed to marshal manifest: %w", err) @@ -118,10 +127,16 @@ func (ms *MediaSignerLocal) SignMP4(ctx context.Context, input io.ReadSeeker, st if err != nil { return nil, fmt.Errorf("failed to unmarshal manifest: %w", err) } + span.End() + + ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_GetSigningAlgorithm") alg, err := c2pa.GetSigningAlgorithm(string(c2pa.ES256K)) if err != nil { return nil, fmt.Errorf("failed to get signing algorithm: %w", err) } + span.End() + + ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_NewBuilder") b, err := c2pa.NewBuilder(&manifest, &c2pa.BuilderParams{ Cert: ms.Cert, Signer: ms.Signer, @@ -131,16 +146,23 @@ func (ms *MediaSignerLocal) SignMP4(ctx context.Context, input io.ReadSeeker, st if err != nil { return nil, fmt.Errorf("failed to create C2PA builder: %w", err) } + span.End() + ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_Sign") output := &aqio.ReadWriteSeeker{} err = b.Sign(input, output, "video/mp4") if err != nil { return nil, fmt.Errorf("failed to sign MP4: %w", err) } + span.End() + + ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_OutputBytes") bs, err := output.Bytes() if err != nil { return nil, fmt.Errorf("failed to get output bytes: %w", err) } + span.End() + spmetrics.SigningDuration.WithLabelValues(ms.StreamerName).Observe(float64(time.Since(startTime).Milliseconds())) return bs, nil } diff --git a/pkg/media/media_signer_ext.go b/pkg/media/media_signer_ext.go index 469c2819f..aff0c10fb 100644 --- a/pkg/media/media_signer_ext.go +++ b/pkg/media/media_signer_ext.go @@ -9,11 +9,14 @@ import ( "io" "os" "os/exec" + "time" "github.com/decred/dcrd/dcrec/secp256k1" "github.com/mr-tron/base58" + "go.opentelemetry.io/otel" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/crypto/aqpub" + "stream.place/streamplace/pkg/spmetrics" ) type MediaSignerExt struct { @@ -52,6 +55,9 @@ func MakeMediaSignerExt(ctx context.Context, cli *config.CLI, streamer string, k } func (ms *MediaSignerExt) SignMP4(ctx context.Context, input io.ReadSeeker, start int64) ([]byte, error) { + startTime := time.Now() + ctx, span := otel.Tracer("signer").Start(ctx, "SignMP4_Ext") + defer span.End() // Get the path to the current executable execPath, err := os.Executable() if err != nil { @@ -98,10 +104,14 @@ func (ms *MediaSignerExt) SignMP4(ctx context.Context, input io.ReadSeeker, star if err := cmd.Wait(); err != nil { return nil, fmt.Errorf("command failed: %w, stderr: %s", err, stderr.String()) } - + spmetrics.SigningDuration.WithLabelValues(ms.streamer).Observe(float64(time.Since(startTime).Milliseconds())) return stdout.Bytes(), nil } func (ms *MediaSignerExt) Pub() aqpub.Pub { return ms.pub } + +func (ms *MediaSignerExt) Streamer() string { + return ms.streamer +} diff --git a/pkg/media/segchanman/segchanman.go b/pkg/media/segchanman/segchanman.go index ae3674042..ba52e9df2 100644 --- a/pkg/media/segchanman/segchanman.go +++ b/pkg/media/segchanman/segchanman.go @@ -4,8 +4,10 @@ import ( "context" "fmt" "sync" + "time" "go.opentelemetry.io/otel" + "stream.place/streamplace/pkg/log" ) // it's a segment channel manager, you see @@ -39,7 +41,7 @@ func (s *SegChanMan) SubscribeSegment(ctx context.Context, user string, renditio chs = []chan *Seg{} s.segChans[key] = chs } - ch := make(chan *Seg, 1024) + ch := make(chan *Seg) chs = append(chs, ch) s.segChans[key] = chs return ch @@ -74,7 +76,14 @@ func (s *SegChanMan) PublishSegment(ctx context.Context, user string, rendition } for _, ch := range chs { go func(ch chan *Seg) { - ch <- seg + select { + case ch <- seg: + case <-ctx.Done(): + return + case <-time.After(1 * time.Minute): + log.Warn(ctx, "failed to send segment to channel, timing out", "user", user, "rendition", rendition) + } + }(ch) } } diff --git a/pkg/media/segmenter.go b/pkg/media/segmenter.go index a7ca52d2d..18e2a3bb0 100644 --- a/pkg/media/segmenter.go +++ b/pkg/media/segmenter.go @@ -9,6 +9,8 @@ import ( "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" "stream.place/streamplace/pkg/log" ) @@ -62,7 +64,9 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) appsink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: WriterNewSample(ctx, buf), EOSFunc: func(sink *app.Sink) { - ctx, span := otel.Tracer("signer").Start(ctx, "SegmentAndSignElem") + ctx, span := otel.Tracer("signer").Start(ctx, "SegmentAndSignElem", trace.WithAttributes( + attribute.String("streamer", ms.Streamer()), + )) defer span.End() resetTimer <- struct{}{} now := time.Now().UnixMilli() diff --git a/pkg/media/webrtc_playback.go b/pkg/media/webrtc_playback.go index e15d8a775..6cd3a3d7c 100644 --- a/pkg/media/webrtc_playback.go +++ b/pkg/media/webrtc_playback.go @@ -49,12 +49,12 @@ func (mm *MediaManager) WebRTCPlayback(ctx context.Context, user string, renditi segBuffer := make(chan *segchanman.Seg, 1024) go func() { + ch := mm.SubscribeSegment(ctx, user, rendition) + defer mm.UnsubscribeSegment(ctx, user, rendition, ch) for { - ch := mm.SubscribeSegment(ctx, user, rendition) select { case <-ctx.Done(): log.Debug(ctx, "exiting segment reader") - mm.UnsubscribeSegment(ctx, user, rendition, ch) return case file := <-ch: log.Debug(ctx, "got segment", "file", file.Filepath) diff --git a/pkg/spmetrics/spmetrics.go b/pkg/spmetrics/spmetrics.go index b0d233be4..b13a75e37 100644 --- a/pkg/spmetrics/spmetrics.go +++ b/pkg/spmetrics/spmetrics.go @@ -48,6 +48,12 @@ var TranscodeDuration = promauto.NewHistogramVec(prometheus.HistogramOpts{ Buckets: []float64{0, 250, 500, 750, 1000, 1250, 1500, 2000, 2500, 3000, 3500, 4000, 4500, 5000, 10000}, }, []string{"streamer"}) +var SigningDuration = promauto.NewHistogramVec(prometheus.HistogramOpts{ + Name: "streamplace_signing_duration_ms", + Help: "duration of transcode in ms", + Buckets: []float64{0, 250, 500, 750, 1000, 1250, 1500, 2000, 2500, 3000, 3500, 4000, 4500, 5000, 10000, 20000, 30000, 60000}, +}, []string{"streamer"}) + var QueuedTranscodeDuration = promauto.NewGaugeVec(prometheus.GaugeOpts{ Name: "streamplace_queued_transcode_duration_ms", Help: "duration of transcode in ms, including time spent waiting", -- 2.51.2