From b3c79b94b997f26a6e481d667f8a97fa1478f8e4 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 10 Jun 2026 14:27:48 -0700 Subject: [PATCH] viewers: count live HLS playback sessions Live HLS playback never touched the viewer count: only the WebRTC path called IncrementViewerCount, so an HLS-only viewer showed as zero. Track sessions by the sid already threaded through every playlist/segment URL: the first request for a (streamer, sid) counts the viewer, and a session idle past 30s uncounts it via a janitor sweep. Media-playlist and segment requests are the heartbeat; master playlist fetches don't count, so one-shot fetchers (preview cards, health checks) aren't viewers. Co-Authored-By: Claude Fable 5 --- pkg/media/hls_sessions.go | 113 ++++++++++++++++++++ pkg/media/hls_sessions_test.go | 106 ++++++++++++++++++ pkg/media/media.go | 13 ++- pkg/spxrpc/place_stream_playback_getlive.go | 9 ++ 4 files changed, 239 insertions(+), 2 deletions(-) create mode 100644 pkg/media/hls_sessions.go create mode 100644 pkg/media/hls_sessions_test.go diff --git a/pkg/media/hls_sessions.go b/pkg/media/hls_sessions.go new file mode 100644 index 000000000..97be32729 --- /dev/null +++ b/pkg/media/hls_sessions.go @@ -0,0 +1,113 @@ +package media + +import ( + "sync" + "time" +) + +// hlsSessionTTL is how long an HLS playback session stays counted after its +// last request. A live player re-fetches the media playlist every target +// duration (~2s) and segments continuously, so a session silent this long is +// gone. Generous enough to ride out a stall/refresh without dropping to zero. +const hlsSessionTTL = 30 * time.Second + +// hlsSessionSweepInterval is how often the janitor checks for idle sessions. +const hlsSessionSweepInterval = 5 * time.Second + +type hlsSessionKey struct { + streamer string + sid string +} + +// hlsSessionTracker turns stateless HLS requests into viewer sessions: the +// first request carrying a given (streamer, sid) starts a session (onNew), +// and a session that goes idle for ttl ends it (onExpire). The two callbacks +// fire exactly once per session, in order, so they can safely drive an +// increment/decrement pair. +type hlsSessionTracker struct { + ttl time.Duration + onNew func(streamer string) + onExpire func(streamer string) + now func() time.Time // clock, overridable in tests + + mu sync.Mutex + lastSeen map[hlsSessionKey]time.Time + janitorRunning bool +} + +func newHLSSessionTracker(ttl time.Duration, onNew, onExpire func(streamer string)) *hlsSessionTracker { + return &hlsSessionTracker{ + ttl: ttl, + onNew: onNew, + onExpire: onExpire, + now: time.Now, + lastSeen: map[hlsSessionKey]time.Time{}, + } +} + +// Touch records a request for (streamer, sid), starting a session if it's +// new. Empty sids are ignored — there's no session to attribute the request +// to. +func (t *hlsSessionTracker) Touch(streamer, sid string) { + if streamer == "" || sid == "" { + return + } + k := hlsSessionKey{streamer: streamer, sid: sid} + t.mu.Lock() + _, known := t.lastSeen[k] + t.lastSeen[k] = t.now() + startJanitor := false + if !known && !t.janitorRunning { + t.janitorRunning = true + startJanitor = true + } + t.mu.Unlock() + if !known { + t.onNew(streamer) + } + if startJanitor { + go t.janitor() + } +} + +// janitor sweeps periodically while any session exists, and exits once the +// tracker drains (the next Touch starts a fresh one). +func (t *hlsSessionTracker) janitor() { + ticker := time.NewTicker(hlsSessionSweepInterval) + defer ticker.Stop() + for range ticker.C { + if t.sweep() { + return + } + } +} + +// sweep expires sessions idle past the ttl and reports whether the tracker is +// now empty (meaning the calling janitor should exit). +func (t *hlsSessionTracker) sweep() bool { + cutoff := t.now().Add(-t.ttl) + t.mu.Lock() + var expired []string + for k, last := range t.lastSeen { + if last.Before(cutoff) { + delete(t.lastSeen, k) + expired = append(expired, k.streamer) + } + } + empty := len(t.lastSeen) == 0 + if empty { + t.janitorRunning = false + } + t.mu.Unlock() + for _, streamer := range expired { + t.onExpire(streamer) + } + return empty +} + +// TouchHLSSession marks HLS playback activity for one (streamer, sid) +// session. Called by the live HLS handlers on media-playlist and segment +// requests; the first touch counts the viewer, going idle uncounts them. +func (mm *MediaManager) TouchHLSSession(streamer, sid string) { + mm.hlsSessions.Touch(streamer, sid) +} diff --git a/pkg/media/hls_sessions_test.go b/pkg/media/hls_sessions_test.go new file mode 100644 index 000000000..549b56ab7 --- /dev/null +++ b/pkg/media/hls_sessions_test.go @@ -0,0 +1,106 @@ +package media + +import ( + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +// countingTracker wires a tracker to a per-streamer net counter so tests can +// assert the onNew/onExpire pairing. +type countingTracker struct { + *hlsSessionTracker + mu sync.Mutex + counts map[string]int + clock time.Time +} + +func newCountingTracker(t *testing.T) *countingTracker { + ct := &countingTracker{ + counts: map[string]int{}, + clock: time.Unix(1000, 0), + } + ct.hlsSessionTracker = newHLSSessionTracker(hlsSessionTTL, + func(streamer string) { + ct.mu.Lock() + defer ct.mu.Unlock() + ct.counts[streamer]++ + }, + func(streamer string) { + ct.mu.Lock() + defer ct.mu.Unlock() + ct.counts[streamer]-- + }, + ) + ct.now = func() time.Time { + ct.mu.Lock() + defer ct.mu.Unlock() + return ct.clock + } + return ct +} + +func (ct *countingTracker) count(streamer string) int { + ct.mu.Lock() + defer ct.mu.Unlock() + return ct.counts[streamer] +} + +func (ct *countingTracker) advance(d time.Duration) { + ct.mu.Lock() + ct.clock = ct.clock.Add(d) + ct.mu.Unlock() +} + +func TestHLSSessionTrackerCountsDistinctSids(t *testing.T) { + ct := newCountingTracker(t) + + ct.Touch("did:plc:streamer", "sid-1") + require.Equal(t, 1, ct.count("did:plc:streamer")) + + // Repeat requests for the same session don't re-count. + ct.Touch("did:plc:streamer", "sid-1") + ct.Touch("did:plc:streamer", "sid-1") + require.Equal(t, 1, ct.count("did:plc:streamer")) + + // A second session counts; sessions are per-streamer. + ct.Touch("did:plc:streamer", "sid-2") + ct.Touch("did:plc:other", "sid-1") + require.Equal(t, 2, ct.count("did:plc:streamer")) + require.Equal(t, 1, ct.count("did:plc:other")) + + // Requests with no sid are unattributable and ignored. + ct.Touch("did:plc:streamer", "") + ct.Touch("", "sid-3") + require.Equal(t, 2, ct.count("did:plc:streamer")) +} + +func TestHLSSessionTrackerExpiry(t *testing.T) { + ct := newCountingTracker(t) + + ct.Touch("did:plc:streamer", "sid-1") + ct.Touch("did:plc:streamer", "sid-2") + require.Equal(t, 2, ct.count("did:plc:streamer")) + + // Within the TTL nothing expires. + ct.advance(hlsSessionTTL / 2) + require.False(t, ct.sweep()) + require.Equal(t, 2, ct.count("did:plc:streamer")) + + // Keep one session alive past the other's TTL. + ct.Touch("did:plc:streamer", "sid-2") + ct.advance(hlsSessionTTL/2 + time.Second) + require.False(t, ct.sweep()) + require.Equal(t, 1, ct.count("did:plc:streamer")) + + // The survivor ages out too; sweep reports the tracker drained. + ct.advance(hlsSessionTTL) + require.True(t, ct.sweep()) + require.Equal(t, 0, ct.count("did:plc:streamer")) + + // A returning sid is a fresh session. + ct.Touch("did:plc:streamer", "sid-1") + require.Equal(t, 1, ct.count("did:plc:streamer")) +} diff --git a/pkg/media/media.go b/pkg/media/media.go index 453598171..7d8a10251 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -54,6 +54,10 @@ type MediaManager struct { webrtcConfig webrtc.Configuration localDB localdb.LocalDB + // hlsSessions tracks live HLS playback sessions by (streamer, sid) so + // stateless HLS requests feed the viewer count. See hls_sessions.go. + hlsSessions *hlsSessionTracker + // Node S2PA transcode signer (cert + PKCS#8 key PEM), built once from the // server-repo key. Used to sign transcode-completed audio tracks under the // node's own did:web identity, signed in-wasm (the node key is software). @@ -113,7 +117,7 @@ func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer if err != nil { return nil, err } - return &MediaManager{ + mm := &MediaManager{ cli: cli, liveWindows: map[string]*livehls.Writer{}, httpPipes: map[string]io.Writer{}, @@ -124,7 +128,12 @@ func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer webrtcConfig: config, localDB: ldb, transcoders: map[string]*streamTranscoder{}, - }, nil + } + mm.hlsSessions = newHLSSessionTracker(hlsSessionTTL, + func(streamer string) { mm.IncrementViewerCount(streamer, "hls") }, + func(streamer string) { mm.DecrementViewerCount(streamer, "hls") }, + ) + return mm, nil } func (mm *MediaManager) HandleData(node *irohStreamplace.PublicKey, data []byte) { diff --git a/pkg/spxrpc/place_stream_playback_getlive.go b/pkg/spxrpc/place_stream_playback_getlive.go index a7190abf0..f1e116ea6 100644 --- a/pkg/spxrpc/place_stream_playback_getlive.go +++ b/pkg/spxrpc/place_stream_playback_getlive.go @@ -98,6 +98,11 @@ func (s *Server) HandleGetLivePlaylist(c echo.Context) error { if body == "" { return echo.NewHTTPError(http.StatusNotFound, "TrackNotFound") } + // A live player re-fetches the media playlist every target duration, + // so this is the session's heartbeat into the viewer count. Master + // requests don't count — one-shot fetchers (preview cards, health + // checks) aren't viewers. + s.mm.TouchHLSSession(did, sid) } h := c.Response().Header() @@ -136,6 +141,10 @@ func (s *Server) HandleGetLiveSegment(c echo.Context) error { return echo.NewHTTPError(http.StatusNotFound, "StreamNotLive") } + // Segment fetches keep the playback session alive in the viewer count + // (sid is threaded through every segment URL by the playlist handler). + s.mm.TouchHLSSession(did, c.QueryParam("sid")) + var data []byte isInit := seg == "init" if isInit { -- 2.51.2