refactor: simplify codebase
This commit is contained in:
+14
-37
@@ -52,12 +52,11 @@ const (
|
||||
|
||||
// StatusEntry is the current miner state for one resolved streamer login.
|
||||
type StatusEntry struct {
|
||||
Login string
|
||||
Watching bool
|
||||
Reason WatchReason
|
||||
WatchedMinutes int
|
||||
WatchStreakMinutes int
|
||||
WatchStreak *int
|
||||
Login string
|
||||
Watching bool
|
||||
Reason WatchReason
|
||||
WatchedMinutes int
|
||||
WatchStreak *int
|
||||
}
|
||||
|
||||
// Manager owns background channel points mining state for the current authenticated user.
|
||||
@@ -89,7 +88,6 @@ type streamerState struct {
|
||||
ChannelPoints int
|
||||
SpadeURL string
|
||||
Metadata *twitch.StreamMetadata
|
||||
LastWatchAt time.Time
|
||||
|
||||
WatchStreak *int
|
||||
WatchStreakMissing bool
|
||||
@@ -98,10 +96,7 @@ type streamerState struct {
|
||||
CurrentWatchReason WatchReason
|
||||
OnlineAt time.Time
|
||||
OfflineAt time.Time
|
||||
PendingStreamUpAt time.Time
|
||||
CurrentBroadcastID string
|
||||
|
||||
seeded bool
|
||||
seeding bool
|
||||
refreshing bool
|
||||
}
|
||||
@@ -146,13 +141,9 @@ func (m *Manager) Sync(ctx context.Context, state *auth.State, viewer *twitch.Vi
|
||||
return
|
||||
}
|
||||
|
||||
type seedTarget struct {
|
||||
configLogin string
|
||||
}
|
||||
|
||||
var (
|
||||
channelIDs []string
|
||||
seedList []seedTarget
|
||||
seedList []string
|
||||
statuses []StatusEntry
|
||||
)
|
||||
|
||||
@@ -176,7 +167,7 @@ func (m *Manager) Sync(ctx context.Context, state *auth.State, viewer *twitch.Vi
|
||||
Login: entry.Login,
|
||||
ChannelID: entry.ChannelID,
|
||||
}
|
||||
seedList = append(seedList, seedTarget{configLogin: entry.ConfigLogin})
|
||||
seedList = append(seedList, entry.ConfigLogin)
|
||||
}
|
||||
|
||||
current.ConfigLogin = entry.ConfigLogin
|
||||
@@ -203,8 +194,8 @@ func (m *Manager) Sync(ctx context.Context, state *auth.State, viewer *twitch.Vi
|
||||
m.emitStatuses(statuses)
|
||||
|
||||
_ = m.pubsub.Sync(ctx, viewer.ID, state.AccessToken, channelIDs)
|
||||
for _, target := range seedList {
|
||||
m.scheduleSeed(target.configLogin)
|
||||
for _, configLogin := range seedList {
|
||||
m.scheduleSeed(configLogin)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -256,7 +247,6 @@ func (m *Manager) scheduleSeed(configLogin string) {
|
||||
m.mu.Lock()
|
||||
if current, ok := m.entries[configLogin]; ok && current.ChannelID == channelID {
|
||||
current.ChannelPoints = channelPoints.Balance
|
||||
current.seeded = true
|
||||
}
|
||||
m.mu.Unlock()
|
||||
m.logf(login, "seeded channel points balance: %d", channelPoints.Balance)
|
||||
@@ -312,7 +302,6 @@ func (m *Manager) scheduleRefresh(configLogin string) {
|
||||
if current, ok := m.entries[configLogin]; ok && current.ChannelID == channelID {
|
||||
current.SpadeURL = spadeURL
|
||||
current.Metadata = metadata
|
||||
current.CurrentBroadcastID = metadata.BroadcastID
|
||||
current.Live = true
|
||||
}
|
||||
m.mu.Unlock()
|
||||
@@ -378,7 +367,6 @@ func (m *Manager) watchOnce(ctx context.Context) error {
|
||||
m.mu.Lock()
|
||||
statuses := []StatusEntry{}
|
||||
if state, ok := m.entries[candidate.configLogin]; ok && state.ChannelID == candidate.channelID {
|
||||
state.LastWatchAt = m.now()
|
||||
state.Watched += time.Minute
|
||||
if candidate.reason == WatchReasonStreak && state.WatchStreakMissing {
|
||||
state.WatchStreakWatched += time.Minute
|
||||
@@ -475,7 +463,6 @@ func (m *Manager) updateWatchReasonsLocked(candidates []watchCandidate) []Status
|
||||
|
||||
func (m *Manager) markOnlineConfirmedLocked(state *streamerState, now time.Time, watchStreak *int) {
|
||||
state.Live = true
|
||||
state.PendingStreamUpAt = time.Time{}
|
||||
if watchStreak != nil {
|
||||
state.WatchStreak = cloneInt(watchStreak)
|
||||
}
|
||||
@@ -487,7 +474,6 @@ func (m *Manager) markOnlineConfirmedLocked(state *streamerState, now time.Time,
|
||||
state.WatchStreakWatched = 0
|
||||
state.SpadeURL = ""
|
||||
state.Metadata = nil
|
||||
state.CurrentBroadcastID = ""
|
||||
}
|
||||
|
||||
func (m *Manager) markOfflineLocked(state *streamerState, now time.Time) {
|
||||
@@ -497,8 +483,6 @@ func (m *Manager) markOfflineLocked(state *streamerState, now time.Time) {
|
||||
state.Metadata = nil
|
||||
state.Watched = 0
|
||||
state.CurrentWatchReason = ""
|
||||
state.PendingStreamUpAt = time.Time{}
|
||||
state.CurrentBroadcastID = ""
|
||||
}
|
||||
|
||||
func (m *Manager) emitStatusForConfig(configLogin string) {
|
||||
@@ -534,12 +518,11 @@ func (m *Manager) emitStatuses(statuses []StatusEntry) {
|
||||
|
||||
func statusFromState(state streamerState) StatusEntry {
|
||||
return StatusEntry{
|
||||
Login: state.Login,
|
||||
Watching: state.Live && state.CurrentWatchReason != "",
|
||||
Reason: state.CurrentWatchReason,
|
||||
WatchedMinutes: int(state.Watched / time.Minute),
|
||||
WatchStreakMinutes: int(state.WatchStreakWatched / time.Minute),
|
||||
WatchStreak: cloneInt(state.WatchStreak),
|
||||
Login: state.Login,
|
||||
Watching: state.Live && state.CurrentWatchReason != "",
|
||||
Reason: state.CurrentWatchReason,
|
||||
WatchedMinutes: int(state.Watched / time.Minute),
|
||||
WatchStreak: cloneInt(state.WatchStreak),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -583,12 +566,6 @@ func (m *Manager) handlePubSubEvent(event Event) {
|
||||
m.logf(state.Login, "claim failed: %v", err)
|
||||
}
|
||||
}
|
||||
case "stream-up":
|
||||
m.mu.Lock()
|
||||
if current, ok := m.entries[configLogin]; ok {
|
||||
current.PendingStreamUpAt = m.now()
|
||||
}
|
||||
m.mu.Unlock()
|
||||
case "stream-down":
|
||||
var statuses []StatusEntry
|
||||
m.mu.Lock()
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
"reflect"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -12,6 +13,7 @@ import (
|
||||
)
|
||||
|
||||
type fakeService struct {
|
||||
mu sync.Mutex
|
||||
channelPoints map[string]*twitch.ChannelPointsContext
|
||||
metadata map[string]*twitch.StreamMetadata
|
||||
spadeURLs map[string]string
|
||||
@@ -20,9 +22,8 @@ type fakeService struct {
|
||||
minuteErr map[string]error
|
||||
claimErr map[string]error
|
||||
|
||||
claimed []string
|
||||
watched []string
|
||||
refreshed []string
|
||||
claimed []string
|
||||
watched []string
|
||||
}
|
||||
|
||||
func (f *fakeService) LoadChannelPointsContext(_ context.Context, login string) (*twitch.ChannelPointsContext, error) {
|
||||
@@ -43,6 +44,8 @@ func (f *fakeService) WatchStreak(_ context.Context, login string) (*int, error)
|
||||
}
|
||||
|
||||
func (f *fakeService) ClaimCommunityPoints(_ context.Context, channelID, claimID string) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.claimed = append(f.claimed, channelID+":"+claimID)
|
||||
if err, ok := f.claimErr[channelID+":"+claimID]; ok {
|
||||
return err
|
||||
@@ -51,7 +54,6 @@ func (f *fakeService) ClaimCommunityPoints(_ context.Context, channelID, claimID
|
||||
}
|
||||
|
||||
func (f *fakeService) StreamMetadata(_ context.Context, login string) (*twitch.StreamMetadata, error) {
|
||||
f.refreshed = append(f.refreshed, login)
|
||||
if result, ok := f.metadata[login]; ok {
|
||||
return result, nil
|
||||
}
|
||||
@@ -73,11 +75,15 @@ func (f *fakeService) PlaybackAccessToken(_ context.Context, login string) (*twi
|
||||
}
|
||||
|
||||
func (f *fakeService) TouchPlayback(_ context.Context, login string, _ *twitch.PlaybackToken) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.watched = append(f.watched, "touch:"+login)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeService) SendMinuteWatched(_ context.Context, spadeURL string, _ twitch.MinuteWatchedPayload) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.watched = append(f.watched, "post:"+spadeURL)
|
||||
if err, ok := f.minuteErr[spadeURL]; ok {
|
||||
return err
|
||||
@@ -85,6 +91,18 @@ func (f *fakeService) SendMinuteWatched(_ context.Context, spadeURL string, _ tw
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeService) claimedCalls() []string {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return append([]string(nil), f.claimed...)
|
||||
}
|
||||
|
||||
func (f *fakeService) watchedCalls() []string {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return append([]string(nil), f.watched...)
|
||||
}
|
||||
|
||||
type fakePubSub struct {
|
||||
syncCalls [][]string
|
||||
}
|
||||
@@ -120,7 +138,7 @@ func TestManagerSyncSeedsAndClaims(t *testing.T) {
|
||||
waitFor(t, func() bool {
|
||||
manager.mu.Lock()
|
||||
defer manager.mu.Unlock()
|
||||
return len(service.claimed) == 1 && manager.entries["alpha"].ChannelPoints == 250 && manager.entries["alpha"].SpadeURL != ""
|
||||
return len(service.claimedCalls()) == 1 && manager.entries["alpha"].ChannelPoints == 250 && manager.entries["alpha"].SpadeURL != ""
|
||||
})
|
||||
|
||||
if got, want := pubsub.syncCalls, [][]string{{"1"}}; !reflect.DeepEqual(got, want) {
|
||||
@@ -146,7 +164,7 @@ func TestManagerWatchOnceUsesTopTwoLiveStreamers(t *testing.T) {
|
||||
}
|
||||
manager := NewManager(context.Background(), service, &fakePubSub{}, nil)
|
||||
defer manager.Close()
|
||||
manager.sleep = cancelSleep
|
||||
manager.watchStarted = true
|
||||
|
||||
manager.Sync(context.Background(), &auth.State{AccessToken: "token"}, &twitch.Viewer{ID: "viewer"}, []twitch.StreamerEntry{
|
||||
{ConfigLogin: "alpha", Login: "alpha_live", ChannelID: "1", Live: true, Status: twitch.StreamerReady},
|
||||
@@ -163,7 +181,7 @@ func TestManagerWatchOnceUsesTopTwoLiveStreamers(t *testing.T) {
|
||||
if err := manager.watchOnce(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got, want := service.watched, []string{
|
||||
if got, want := service.watchedCalls(), []string{
|
||||
"touch:alpha_live",
|
||||
"post:https://spade.test/alpha",
|
||||
"touch:beta_live",
|
||||
@@ -189,7 +207,7 @@ func TestManagerWatchOncePrioritizesMissingWatchStreaks(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if got, want := service.watched, []string{
|
||||
if got, want := service.watchedCalls(), []string{
|
||||
"touch:beta_live",
|
||||
"post:https://spade.test/beta_live",
|
||||
"touch:gamma_live",
|
||||
@@ -215,7 +233,7 @@ func TestManagerWatchOnceFillsWithPointsCandidates(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if got, want := service.watched, []string{
|
||||
if got, want := service.watchedCalls(), []string{
|
||||
"touch:beta_live",
|
||||
"post:https://spade.test/beta_live",
|
||||
"touch:alpha_live",
|
||||
@@ -335,7 +353,7 @@ func TestHandlePubSubWatchStreakCompletesMaintenance(t *testing.T) {
|
||||
CurrentWatchReason: WatchReasonStreak,
|
||||
}
|
||||
|
||||
manager.handlePubSubEvent(Event{MessageType: "points-earned", ChannelID: "1", Balance: 555, ReasonCode: "WATCH_STREAK", TotalPoints: 450, Timestamp: "t1", Topic: "community-points-user-v1"})
|
||||
manager.handlePubSubEvent(Event{MessageType: "points-earned", ChannelID: "1", Balance: 555, ReasonCode: "WATCH_STREAK", Timestamp: "t1", Topic: "community-points-user-v1"})
|
||||
|
||||
waitFor(t, func() bool {
|
||||
manager.mu.Lock()
|
||||
@@ -434,26 +452,6 @@ func TestManagerSyncPollingConfirmedOnlineResetsStreakState(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestManagerStreamUpWaitsForConfirmation(t *testing.T) {
|
||||
now := time.Date(2026, 4, 30, 12, 0, 0, 0, time.UTC)
|
||||
manager := NewManager(context.Background(), &fakeService{}, &fakePubSub{}, nil)
|
||||
defer manager.Close()
|
||||
manager.now = func() time.Time { return now }
|
||||
manager.entries["alpha"] = &streamerState{ConfigLogin: "alpha", ChannelID: "1", Login: "alpha_live"}
|
||||
|
||||
manager.handlePubSubEvent(Event{MessageType: "stream-up", ChannelID: "1", Timestamp: "up", Topic: "video-playback-by-id"})
|
||||
|
||||
manager.mu.Lock()
|
||||
defer manager.mu.Unlock()
|
||||
state := manager.entries["alpha"]
|
||||
if state.Live {
|
||||
t.Fatal("stream-up should wait for API/viewcount confirmation")
|
||||
}
|
||||
if !state.PendingStreamUpAt.Equal(now) {
|
||||
t.Fatalf("pending stream-up = %s, want %s", state.PendingStreamUpAt, now)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandlePubSubEventUpdatesBalanceAndClaims(t *testing.T) {
|
||||
service := &fakeService{}
|
||||
manager := NewManager(context.Background(), service, &fakePubSub{}, nil)
|
||||
@@ -468,7 +466,7 @@ func TestHandlePubSubEventUpdatesBalanceAndClaims(t *testing.T) {
|
||||
if manager.entries["alpha"].ChannelPoints != 555 {
|
||||
t.Fatalf("channelPoints = %d", manager.entries["alpha"].ChannelPoints)
|
||||
}
|
||||
if got, want := service.claimed, []string{"1:claim-2"}; !reflect.DeepEqual(got, want) {
|
||||
if got, want := service.claimedCalls(), []string{"1:claim-2"}; !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("claimed = %#v, want %#v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
+15
-46
@@ -6,7 +6,8 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sort"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -29,7 +30,6 @@ type Event struct {
|
||||
Balance int
|
||||
ClaimID string
|
||||
ReasonCode string
|
||||
TotalPoints int
|
||||
}
|
||||
|
||||
func (e Event) key() string {
|
||||
@@ -45,7 +45,6 @@ type Client struct {
|
||||
dialer *websocket.Dialer
|
||||
mu sync.Mutex
|
||||
conn *websocket.Conn
|
||||
viewerID string
|
||||
token string
|
||||
topics []string
|
||||
updateCh chan struct{}
|
||||
@@ -81,12 +80,11 @@ func (c *Client) Sync(_ context.Context, viewerID, accessToken string, channelID
|
||||
}
|
||||
topics = append(topics, "video-playback-by-id."+channelID)
|
||||
}
|
||||
sort.Strings(topics)
|
||||
slices.Sort(topics)
|
||||
|
||||
c.mu.Lock()
|
||||
c.viewerID = viewerID
|
||||
c.token = accessToken
|
||||
changed := !equalStrings(c.topics, topics)
|
||||
changed := !slices.Equal(c.topics, topics)
|
||||
c.topics = topics
|
||||
conn := c.conn
|
||||
c.mu.Unlock()
|
||||
@@ -258,15 +256,13 @@ func (c *Client) snapshot() ([]string, string) {
|
||||
|
||||
type frame struct {
|
||||
Type string
|
||||
Error string
|
||||
Event Event
|
||||
}
|
||||
|
||||
func parseFrame(message []byte) (*frame, error) {
|
||||
var envelope struct {
|
||||
Type string `json:"type"`
|
||||
Error string `json:"error"`
|
||||
Data *struct {
|
||||
Type string `json:"type"`
|
||||
Data *struct {
|
||||
Topic string `json:"topic"`
|
||||
Message string `json:"message"`
|
||||
} `json:"data"`
|
||||
@@ -275,7 +271,7 @@ func parseFrame(message []byte) (*frame, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
result := &frame{Type: envelope.Type, Error: envelope.Error}
|
||||
result := &frame{Type: envelope.Type}
|
||||
if envelope.Type != "MESSAGE" || envelope.Data == nil {
|
||||
return result, nil
|
||||
}
|
||||
@@ -294,14 +290,18 @@ func parseFrame(message []byte) (*frame, error) {
|
||||
ChannelID string `json:"channel_id"`
|
||||
} `json:"balance"`
|
||||
PointGain *struct {
|
||||
ReasonCode string `json:"reason_code"`
|
||||
TotalPoints int `json:"total_points"`
|
||||
ReasonCode string `json:"reason_code"`
|
||||
} `json:"point_gain"`
|
||||
} `json:"data"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(envelope.Data.Message), &payload); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
topic := envelope.Data.Topic
|
||||
topicName, topicChannelID := topic, ""
|
||||
if i := strings.LastIndexByte(topic, '.'); i >= 0 {
|
||||
topicName, topicChannelID = topic[:i], topic[i+1:]
|
||||
}
|
||||
|
||||
channelID := payload.Data.ChannelID
|
||||
if payload.Data.Claim != nil && payload.Data.Claim.ChannelID != "" {
|
||||
@@ -311,11 +311,11 @@ func parseFrame(message []byte) (*frame, error) {
|
||||
channelID = payload.Data.Balance.ChannelID
|
||||
}
|
||||
if channelID == "" {
|
||||
channelID = topicSuffix(envelope.Data.Topic)
|
||||
channelID = topicChannelID
|
||||
}
|
||||
|
||||
result.Event = Event{
|
||||
Topic: topicPrefix(envelope.Data.Topic),
|
||||
Topic: topicName,
|
||||
MessageType: payload.Type,
|
||||
ChannelID: channelID,
|
||||
Timestamp: payload.Data.Timestamp,
|
||||
@@ -328,37 +328,6 @@ func parseFrame(message []byte) (*frame, error) {
|
||||
}
|
||||
if payload.Data.PointGain != nil {
|
||||
result.Event.ReasonCode = payload.Data.PointGain.ReasonCode
|
||||
result.Event.TotalPoints = payload.Data.PointGain.TotalPoints
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func topicPrefix(topic string) string {
|
||||
for i := len(topic) - 1; i >= 0; i-- {
|
||||
if topic[i] == '.' {
|
||||
return topic[:i]
|
||||
}
|
||||
}
|
||||
return topic
|
||||
}
|
||||
|
||||
func topicSuffix(topic string) string {
|
||||
for i := len(topic) - 1; i >= 0; i-- {
|
||||
if topic[i] == '.' {
|
||||
return topic[i+1:]
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func equalStrings(left, right []string) bool {
|
||||
if len(left) != len(right) {
|
||||
return false
|
||||
}
|
||||
for i := range left {
|
||||
if left[i] != right[i] {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -21,7 +21,7 @@ func TestParseFramePointsEarnedPointGain(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if frame.Event.ReasonCode != "WATCH_STREAK" || frame.Event.TotalPoints != 450 {
|
||||
if frame.Event.ReasonCode != "WATCH_STREAK" {
|
||||
t.Fatalf("event = %#v", frame.Event)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user