From b9b4193af347cd7d84aac712a424c0f19a75ea65 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Mon, 20 Jul 2026 10:29:53 -0500 Subject: [PATCH] director: packetize+publish playback segments in order per rendition MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit NewSegment dispatched AddPlaybackSegment through ss.Go, which dispatched AddToWebRTC through ss.Go again — so each segment's packetize (a full gst pipeline, ~100-400ms with high variance) raced every other segment's and PublishSegment order could invert. On streams with short GoPs (1s keyint), packetize latency variance is comparable to the inter-segment interval, so inversions fire constantly; WebRTC playback, the only consumer that plays segments strictly in arrival order, showed each swap as frame loss at a keyframe. Each rendition now gets its own ordered queue + worker: publish order always equals enqueue order, one rendition's backlog never delays another's, and a full queue drops the incoming segment (counted in streamplace_playback_queue_dropped_total) rather than blocking the director's segment loop. Workers drain on session teardown so the stream's tail still publishes. Also removes the dead StreamSession.packets field. --- go.mod | 1 + go.sum | 2 - pkg/director/director.go | 1 - pkg/director/stream_session.go | 108 +++++++++++++++++++---- pkg/director/stream_session_test.go | 128 ++++++++++++++++++++++++++++ pkg/spmetrics/spmetrics.go | 5 ++ 6 files changed, 224 insertions(+), 21 deletions(-) diff --git a/go.mod b/go.mod index 16a3cf5ae..16c7ae87d 100644 --- a/go.mod +++ b/go.mod @@ -364,6 +364,7 @@ require ( github.com/kolesa-team/go-webp v1.0.5 // indirect github.com/kulti/thelper v0.6.3 // indirect github.com/kunwardeep/paralleltest v1.0.14 // indirect + github.com/kylelemons/godebug v1.1.0 // indirect github.com/labstack/gommon v0.4.2 // indirect github.com/lasiar/canonicalheader v1.1.2 // indirect github.com/ldez/exptostd v0.4.3 // indirect diff --git a/go.sum b/go.sum index 7f1411407..ee7c1b6e7 100644 --- a/go.sum +++ b/go.sum @@ -1372,8 +1372,6 @@ github.com/streamplace/atmoq/go v0.0.4-0.20260701223355-13757de4ae08 h1:NiTRz8AX github.com/streamplace/atmoq/go v0.0.4-0.20260701223355-13757de4ae08/go.mod h1:3P8eSwKAGH7uh3SX5z1jlt/JgPTilJTUZngQJKhWY5s= github.com/streamplace/atproto-oauth-golang v0.0.0-20260413212710-98956064d06c h1:IzEPU2O4iL58Nb7aw+7lB9ttnesEwOVVE5oV9NEXemM= github.com/streamplace/atproto-oauth-golang v0.0.0-20260413212710-98956064d06c/go.mod h1:9LlKkqciiO5lRfbX0n4Wn5KNY9nvFb4R3by8FdW2TWc= -github.com/streamplace/glex v0.0.0-20260715231618-ee553e32d7c7 h1:MSBBIH+QMR9AVfC0RuBLbBm/o1RAl7+bacek7txswDo= -github.com/streamplace/glex v0.0.0-20260715231618-ee553e32d7c7/go.mod h1:LRaoeSMvSgOrhFX8s7ygjRlyka7wXdDa1s7JJ9o1IzY= github.com/streamplace/glex v0.0.0-20260716203108-f73ed7cc31c9 h1:HbIhx8i7wytiNg9mWg2EuHG59r8FXTVlDqlkZObjqbQ= github.com/streamplace/glex v0.0.0-20260716203108-f73ed7cc31c9/go.mod h1:LRaoeSMvSgOrhFX8s7ygjRlyka7wXdDa1s7JJ9o1IzY= github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4 h1:L1fS4HJSaAyNnkwfuZubgfeZy8rkWmA0cMtH5Z0HqNc= diff --git a/pkg/director/director.go b/pkg/director/director.go index 82ce16fe6..e16a70aa3 100644 --- a/pkg/director/director.go +++ b/pkg/director/director.go @@ -76,7 +76,6 @@ func (d *Director) Start(ctx context.Context) error { bus: d.bus, segmentChan: make(chan struct{}), op: d.op, - packets: make([]bus.PacketizedSegment, 0), started: make(chan struct{}), statefulDB: d.statefulDB, replicator: d.replicator, diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index cf6605de4..812ef657a 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -8,6 +8,7 @@ import ( "io" "net/url" "strings" + "sync" "time" "github.com/bluesky-social/indigo/atproto/syntax" @@ -63,11 +64,20 @@ type StreamSession struct { g *errgroup.Group started chan struct{} ctx context.Context - packets []bus.PacketizedSegment statefulDB *statedb.StatefulDB replicator replication.Replicator atsync *atproto.ATProtoSynchronizer + // playbackWorkers holds one ordered packetize+publish queue per rendition. + // Packetizing a segment spins up a full gst pipeline (~100-400ms, high + // variance); feeding segments to subscribers from per-segment goroutines + // let segment N+1 publish before segment N, which WebRTC playback — a + // strict arrival-order consumer — showed as frame loss at every keyframe. + playbackWorkers map[string]chan playbackJob + playbackWorkersMu sync.Mutex + // addToWebRTCFn is a test seam; nil means AddToWebRTC. + addToWebRTCFn func(ctx context.Context, spseg *placestream.Segment, rendition string, seg *bus.Seg) error + lastLivestreamTime time.Time lastViewCountTime time.Time s3Uploader *s3.S3Uploader @@ -274,13 +284,11 @@ func (ss *StreamSession) NewSegment(ctx context.Context, notif *media.NewSegment } ss.bus.Publish(spseg.Creator, spseg) - ss.Go(ctx, func() error { - return ss.AddPlaybackSegment(ctx, spseg, "source", &bus.Seg{ - Filepath: notif.Segment.ID, - Data: notif.Data, - Muxl: notif.Muxl, - Published: notif.Metadata.Published, - }) + ss.AddPlaybackSegment(ctx, spseg, "source", &bus.Seg{ + Filepath: notif.Segment.ID, + Data: notif.Data, + Muxl: notif.Muxl, + Published: notif.Metadata.Published, }) if notif.Local { @@ -875,22 +883,86 @@ func (ss *StreamSession) Transcode(ctx context.Context, spseg *placestream.Segme if err != nil { return fmt.Errorf("failed to write transcoded segment file: %w", err) } - ss.Go(ctx, func() error { - return ss.AddPlaybackSegment(ctx, spseg, rs[i].Name, &bus.Seg{ - Filepath: fd.Name(), - Data: seg, - }) + // NOTE: renditions from concurrent Transcode calls can still enqueue + // out of segment order (transcode latency varies per segment); the + // ordered worker only guarantees publish order == enqueue order. + ss.AddPlaybackSegment(ctx, spseg, rs[i].Name, &bus.Seg{ + Filepath: fd.Name(), + Data: seg, }) } return nil } -func (ss *StreamSession) AddPlaybackSegment(ctx context.Context, spseg *placestream.Segment, rendition string, seg *bus.Seg) error { - ss.Go(ctx, func() error { - return ss.AddToWebRTC(ctx, spseg, rendition, seg) - }) - return nil +// playbackJob is one segment awaiting packetization + publish on its +// rendition's ordered queue. ctx is the caller's (director) context — it +// outlives the session's own cancellation so queued segments can still be +// packetized while the worker drains on shutdown. +type playbackJob struct { + ctx context.Context + spseg *placestream.Segment + rendition string + seg *bus.Seg +} + +// AddPlaybackSegment queues a segment for packetization + publish to playback +// subscribers. Each rendition has its own queue and worker, so publish order +// always equals enqueue order and one rendition's backlog never delays +// another's. A full queue (packetize wedged for dozens of segments) drops the +// incoming segment rather than blocking the director's segment loop. +func (ss *StreamSession) AddPlaybackSegment(ctx context.Context, spseg *placestream.Segment, rendition string, seg *bus.Seg) { + ss.playbackWorkersMu.Lock() + if ss.playbackWorkers == nil { + ss.playbackWorkers = map[string]chan playbackJob{} + } + q, ok := ss.playbackWorkers[rendition] + if !ok { + q = make(chan playbackJob, 64) + ss.playbackWorkers[rendition] = q + ss.g.Go(func() error { + return ss.playbackWorker(ss.ctx, rendition, q) + }) + } + ss.playbackWorkersMu.Unlock() + select { + case q <- playbackJob{ctx: ctx, spseg: spseg, rendition: rendition, seg: seg}: + case <-ss.ctx.Done(): + default: + spmetrics.PlaybackQueueDropped.WithLabelValues(spseg.Creator, rendition).Inc() + log.Error(ctx, "playback queue full, dropping segment", "rendition", rendition, "segID", seg.Filepath) + } +} + +// playbackWorker packetizes and publishes one rendition's segments one at a +// time, in queue order. On session teardown it drains the queue before +// exiting so the stream's tail still publishes. +func (ss *StreamSession) playbackWorker(ctx context.Context, rendition string, q chan playbackJob) error { + for { + select { + case <-ctx.Done(): + for { + select { + case job := <-q: + ss.processPlaybackJob(job) + default: + return nil + } + } + case job := <-q: + ss.processPlaybackJob(job) + } + } +} + +func (ss *StreamSession) processPlaybackJob(job playbackJob) { + fn := ss.addToWebRTCFn + if fn == nil { + fn = ss.AddToWebRTC + } + if err := fn(job.ctx, job.spseg, job.rendition, job.seg); err != nil { + log.Error(job.ctx, "error in playback worker", "error", err, "rendition", job.rendition) + } } func (ss *StreamSession) AddToWebRTC(ctx context.Context, spseg *placestream.Segment, rendition string, seg *bus.Seg) error { diff --git a/pkg/director/stream_session_test.go b/pkg/director/stream_session_test.go index 4852606fd..047d173b1 100644 --- a/pkg/director/stream_session_test.go +++ b/pkg/director/stream_session_test.go @@ -1,10 +1,20 @@ package director import ( + "context" + "fmt" + "sync" "testing" "time" "github.com/stretchr/testify/require" + "golang.org/x/sync/errgroup" + "stream.place/streamplace/pkg/bus" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/placestream" + "stream.place/streamplace/pkg/spmetrics" + + "github.com/prometheus/client_golang/prometheus/testutil" ) func TestExceedsMaxBitrate(t *testing.T) { @@ -51,3 +61,121 @@ func TestExceedsMaxBitrateMarginBoundary(t *testing.T) { require.True(t, justOutside, "8Mbit beyond 10%% of 7.2Mbit max should kick") require.Equal(t, eightMbit, func() int { r, _ := exceedsMaxBitrate(megabyte, time.Second.Nanoseconds(), 1); return r }()) } + +func newTestStreamSession(t *testing.T) (*StreamSession, context.CancelFunc) { + t.Helper() + ctx, cancel := context.WithCancel(context.Background()) + g, gctx := errgroup.WithContext(ctx) + started := make(chan struct{}) + close(started) + ss := &StreamSession{ + cli: &config.CLI{}, + bus: bus.NewBus(), + g: g, + ctx: gctx, + started: started, + } + return ss, cancel +} + +// Playback segments must be packetized + published strictly in enqueue order: +// the old ss.Go dispatch ran each segment's packetize concurrently, so a slow +// segment N could publish after segment N+1 — and WebRTC playback, which +// consumes segments strictly in arrival order, showed the swap as frame loss +// at every keyframe. +func TestPlaybackWorkerPublishesInSegmentOrder(t *testing.T) { + ss, cancel := newTestStreamSession(t) + segChan := ss.bus.SubscribeSegment(context.Background(), "did:test:streamer", "source") + ss.addToWebRTCFn = func(ctx context.Context, spseg *placestream.Segment, rendition string, seg *bus.Seg) error { + ss.bus.PublishSegment(ctx, spseg.Creator, rendition, seg) + return nil + } + // under the queue cap so the drop-on-full policy can't fire — this test + // isolates ordering + const n = 50 + for i := 0; i < n; i++ { + ss.AddPlaybackSegment(context.Background(), &placestream.Segment{Creator: "did:test:streamer"}, "source", &bus.Seg{Filepath: fmt.Sprintf("seg-%05d", i)}) + } + cancel() + require.NoError(t, ss.g.Wait()) + for i := 0; i < n; i++ { + select { + case seg := <-segChan.C: + require.Equal(t, fmt.Sprintf("seg-%05d", i), seg.Filepath) + case <-time.After(5 * time.Second): + t.Fatalf("timed out waiting for segment %d", i) + } + } +} + +// A wedged packetize pipeline must not stall the director's segment loop: +// once a rendition's queue is full, new segments for it drop (loudly) instead +// of blocking enqueue. +func TestAddPlaybackSegmentDropsWhenQueueFull(t *testing.T) { + ss, cancel := newTestStreamSession(t) + gate := make(chan struct{}) + entered := make(chan struct{}) + var once sync.Once + processed := make(chan string, 128) + ss.addToWebRTCFn = func(ctx context.Context, spseg *placestream.Segment, rendition string, seg *bus.Seg) error { + once.Do(func() { close(entered) }) + <-gate + processed <- seg.Filepath + return nil + } + enqueue := func(i int) { + ss.AddPlaybackSegment(context.Background(), &placestream.Segment{Creator: "did:test:streamer"}, "source", &bus.Seg{Filepath: fmt.Sprintf("seg-%05d", i)}) + } + enqueue(0) + <-entered // worker now wedged inside packetize; queue empty + const queueCap = 64 + for i := 1; i <= queueCap; i++ { + enqueue(i) + } + enqueue(queueCap + 1) // queue full — this one drops + require.Equal(t, float64(1), testutil.ToFloat64(spmetrics.PlaybackQueueDropped.WithLabelValues("did:test:streamer", "source"))) + close(gate) + cancel() + require.NoError(t, ss.g.Wait()) + got := []string{} + for { + select { + case fp := <-processed: + got = append(got, fp) + default: + require.Len(t, got, queueCap+1) + require.NotContains(t, got, fmt.Sprintf("seg-%05d", queueCap+1)) + return + } + } +} + +// Each rendition gets its own worker: one rendition's wedged queue must not +// delay another rendition's segments. +func TestPlaybackWorkersArePerRendition(t *testing.T) { + ss, cancel := newTestStreamSession(t) + gate := make(chan struct{}) + entered := make(chan struct{}) + var once sync.Once + processed := make(chan string, 16) + ss.addToWebRTCFn = func(ctx context.Context, spseg *placestream.Segment, rendition string, seg *bus.Seg) error { + if rendition == "source" { + once.Do(func() { close(entered) }) + <-gate + } + processed <- rendition + ":" + seg.Filepath + return nil + } + ss.AddPlaybackSegment(context.Background(), &placestream.Segment{Creator: "did:test:streamer"}, "source", &bus.Seg{Filepath: "seg-00000"}) + <-entered // source worker now wedged + ss.AddPlaybackSegment(context.Background(), &placestream.Segment{Creator: "did:test:streamer"}, "other", &bus.Seg{Filepath: "seg-00001"}) + select { + case got := <-processed: + require.Equal(t, "other:seg-00001", got) + case <-time.After(5 * time.Second): + t.Fatal("one rendition's wedged queue blocked another rendition") + } + close(gate) + cancel() + require.NoError(t, ss.g.Wait()) +} diff --git a/pkg/spmetrics/spmetrics.go b/pkg/spmetrics/spmetrics.go index d3a44ea21..5bbb30208 100644 --- a/pkg/spmetrics/spmetrics.go +++ b/pkg/spmetrics/spmetrics.go @@ -103,6 +103,11 @@ var SegmentPublishDropped = promauto.NewCounterVec(prometheus.CounterOpts{ Help: "segments dropped for a subscriber whose queue was a full buffer behind", }, []string{"streamer", "rendition"}) +var PlaybackQueueDropped = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "streamplace_playback_queue_dropped_total", + Help: "playback segments dropped because the session's packetize queue was full", +}, []string{"streamer", "rendition"}) + var LabelerFirehosesConnected = promauto.NewGaugeVec(prometheus.GaugeOpts{ Name: "streamplace_labeler_firehoses_connected", Help: "number of currently connected labeler firehoses", -- 2.51.2