diff --git a/pkg/api/api.go b/pkg/api/api.go index b77aec72..a14fe872 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -41,7 +41,6 @@ import ( "stream.place/streamplace/pkg/mist/mistconfig" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/notifications" - "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/spxrpc" "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/streamplace" @@ -574,7 +573,7 @@ func (a *StreamplaceAPI) HandleViewCount(ctx context.Context) httprouter.Handle apierrors.WriteHTTPNotFound(w, "user not found", err) return } - count := spmetrics.GetViewCount(user) + count := a.Bus.GetViewerCount(user) bs, err := json.Marshal(streamplace.Livestream_ViewerCount{Count: int64(count), LexiconTypeID: "place.stream.livestream#viewerCount"}) if err != nil { apierrors.WriteHTTPInternalServerError(w, "could not marshal view count", err) diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 13610a55..15d7bd1d 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -559,6 +559,10 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD // This allows moderators to see their permissions instantly without page refresh go atsync.Bus.Publish(userDID, view) + case *streamplace.LiveViewCount: + log.Debug(ctx, "indexing view count", "streamer", rec.Streamer, "server", rec.Server, "count", rec.Count) + atsync.Bus.SetFederatedViewCount(rec.Streamer, rec.Server, int(rec.Count)) + case *streamplace.LiveRecommendations: log.Debug(ctx, "creating recommendations", "userDID", userDID, "count", len(rec.Streamers)) diff --git a/pkg/bus/bus.go b/pkg/bus/bus.go index 9f004223..3b8335b7 100644 --- a/pkg/bus/bus.go +++ b/pkg/bus/bus.go @@ -1,7 +1,11 @@ package bus import ( + "strings" "sync" + "time" + + "github.com/patrickmn/go-cache" ) type Message any @@ -24,6 +28,11 @@ type Bus struct { viewerCounts map[string]map[string]int viewerCountsMutex sync.RWMutex viewerCountSubscriptions []chan ViewerCountUpdate + + // federatedViewCounts stores reported view counts from remote servers. + // Key: "{streamerDID}:{serverDID}", Value: int (count). + // 2min TTL so stale servers drop off (records update every ~30s). + federatedViewCounts *cache.Cache } func NewBus() *Bus { @@ -33,6 +42,7 @@ func NewBus() *Bus { segBuf: make(map[string][]*Seg), viewerCounts: make(map[string]map[string]int), viewerCountSubscriptions: []chan ViewerCountUpdate{}, + federatedViewCounts: cache.New(2*time.Minute, 4*time.Minute), } } @@ -95,17 +105,32 @@ func (b *Bus) Publish(user string, msg Message) { func (b *Bus) GetViewerCount(user string) int { b.viewerCountsMutex.RLock() defer b.viewerCountsMutex.RUnlock() - streamerCounts, ok := b.viewerCounts[user] - if !ok { - return 0 - } + // Local viewer counts (from HLS connections on this node) count := 0 - for _, viewers := range streamerCounts { - count += viewers + if streamerCounts, ok := b.viewerCounts[user]; ok { + for _, viewers := range streamerCounts { + count += viewers + } + } + // Federated view counts (reported by remote servers via atproto) + prefix := user + ":" + for k, v := range b.federatedViewCounts.Items() { + if strings.HasPrefix(k, prefix) && !v.Expired() { + if c, ok := v.Object.(int); ok { + count += c + } + } } return count } +// SetFederatedViewCount stores a view count reported by a remote server. +// The count will expire after 2 minutes if not refreshed. +func (b *Bus) SetFederatedViewCount(streamer string, server string, count int) { + key := streamer + ":" + server + b.federatedViewCounts.SetDefault(key, count) +} + func (b *Bus) SetViewerCount(user string, origin string, count int) { b.viewerCountsMutex.Lock() defer b.viewerCountsMutex.Unlock() diff --git a/pkg/spxrpc/place_stream_live.go b/pkg/spxrpc/place_stream_live.go index fcf69183..0caa6cce 100644 --- a/pkg/spxrpc/place_stream_live.go +++ b/pkg/spxrpc/place_stream_live.go @@ -190,7 +190,7 @@ func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before s if err != nil { return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("Failed to convert livestream to streamplace livestream: %s", err)) } - viewers := spmetrics.GetViewCount(stream.Author.Did) + viewers := s.bus.GetViewerCount(stream.Author.Did) stream.ViewerCount = &placestream.Livestream_ViewerCount{ LexiconTypeID: "place.stream.livestream#viewerCount", Count: int64(viewers),