mirror of
https://github.com/navidrome/navidrome.git
synced 2026-08-31 07:30:32 +00:00
* feat(scrobbler): add exponential backoff delay helper * feat(scrobbler): back off retries up to 4m during outages * test(scrobbler): verify backoff schedule with synctest; clarify backoffDelay doc Adds a testing/synctest-based test that drives the real run loop against a failing service and asserts the exact 5s/10s/20s/40s retry schedule and the drain-on-recovery reset, addressing the review note that the run loop's behavior was untested. Also clarifies the backoffDelay doc comment: the argument is a zero-based retry index.
212 lines
5.8 KiB
Go
212 lines
5.8 KiB
Go
package scrobbler
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"time"
|
|
|
|
"github.com/navidrome/navidrome/log"
|
|
"github.com/navidrome/navidrome/model"
|
|
"github.com/navidrome/navidrome/model/request"
|
|
)
|
|
|
|
const (
|
|
minRetryDelay = 5 * time.Second
|
|
maxRetryDelay = 4 * time.Minute
|
|
// maxRetryShift caps the exponent so the shift never overflows int64.
|
|
// minRetryDelay<<6 = 320s already exceeds maxRetryDelay, so 6 reaches the ceiling.
|
|
maxRetryShift = 6
|
|
)
|
|
|
|
// backoffDelay returns the delay for a zero-based retry index (0 = first retry):
|
|
// minRetryDelay doubled per prior failure, clamped to maxRetryDelay.
|
|
func backoffDelay(failures int) time.Duration {
|
|
if failures < 0 {
|
|
failures = 0
|
|
}
|
|
if failures >= maxRetryShift {
|
|
return maxRetryDelay
|
|
}
|
|
d := minRetryDelay << failures
|
|
if d > maxRetryDelay {
|
|
return maxRetryDelay
|
|
}
|
|
return d
|
|
}
|
|
|
|
// Loader is a function that loads a scrobbler by name.
|
|
// It returns the scrobbler and true if found, or nil and false if not available.
|
|
// This allows the buffered scrobbler to always get the current plugin instance.
|
|
type Loader func() (Scrobbler, bool)
|
|
|
|
// newBufferedScrobbler creates a buffered scrobbler that wraps a static scrobbler instance.
|
|
// Use this for builtin scrobblers that don't change.
|
|
func newBufferedScrobbler(ds model.DataStore, s Scrobbler, service string) *bufferedScrobbler {
|
|
return newBufferedScrobblerWithLoader(ds, service, func() (Scrobbler, bool) {
|
|
return s, true
|
|
})
|
|
}
|
|
|
|
// newBufferedScrobblerWithLoader creates a buffered scrobbler that dynamically loads
|
|
// the underlying scrobbler on each call. Use this for plugin scrobblers that may be
|
|
// reloaded (e.g., after configuration changes).
|
|
func newBufferedScrobblerWithLoader(ds model.DataStore, service string, loader Loader) *bufferedScrobbler {
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
b := &bufferedScrobbler{
|
|
ds: ds,
|
|
loader: loader,
|
|
service: service,
|
|
wakeSignal: make(chan struct{}, 1),
|
|
ctx: ctx,
|
|
cancel: cancel,
|
|
}
|
|
go b.run(ctx)
|
|
return b
|
|
}
|
|
|
|
type bufferedScrobbler struct {
|
|
ds model.DataStore
|
|
loader Loader
|
|
service string
|
|
wakeSignal chan struct{}
|
|
ctx context.Context
|
|
cancel context.CancelFunc
|
|
}
|
|
|
|
func (b *bufferedScrobbler) Stop() {
|
|
if b.cancel != nil {
|
|
b.cancel()
|
|
}
|
|
}
|
|
|
|
func (b *bufferedScrobbler) IsAuthorized(ctx context.Context, userId string) bool {
|
|
s, ok := b.loader()
|
|
if !ok {
|
|
return false
|
|
}
|
|
return s.IsAuthorized(ctx, userId)
|
|
}
|
|
|
|
func (b *bufferedScrobbler) NowPlaying(ctx context.Context, userId string, track *model.MediaFile, position int) error {
|
|
s, ok := b.loader()
|
|
if !ok {
|
|
return errors.New("scrobbler not available")
|
|
}
|
|
return s.NowPlaying(ctx, userId, track, position)
|
|
}
|
|
|
|
func (b *bufferedScrobbler) Scrobble(ctx context.Context, userId string, s Scrobble) error {
|
|
err := b.ds.ScrobbleBuffer(ctx).Enqueue(b.service, userId, s.ID, s.TimeStamp)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
b.sendWakeSignal()
|
|
return nil
|
|
}
|
|
|
|
func (b *bufferedScrobbler) PlaybackReport(ctx context.Context, info PlaybackSession) error {
|
|
s, ok := b.loader()
|
|
if !ok {
|
|
return errors.New("scrobbler not available")
|
|
}
|
|
return s.PlaybackReport(ctx, info)
|
|
}
|
|
|
|
func (b *bufferedScrobbler) sendWakeSignal() {
|
|
// Don't block if the previous signal was not read yet
|
|
select {
|
|
case b.wakeSignal <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (b *bufferedScrobbler) run(ctx context.Context) {
|
|
timer := time.NewTimer(time.Hour)
|
|
timer.Stop()
|
|
defer timer.Stop()
|
|
failures := 0
|
|
for {
|
|
if b.processQueue(ctx) {
|
|
failures = 0
|
|
timer.Stop()
|
|
} else {
|
|
timer.Reset(backoffDelay(failures))
|
|
if failures < maxRetryShift {
|
|
failures++
|
|
}
|
|
}
|
|
select {
|
|
case <-b.wakeSignal:
|
|
case <-timer.C:
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func (b *bufferedScrobbler) processQueue(ctx context.Context) bool {
|
|
buffer := b.ds.ScrobbleBuffer(ctx)
|
|
userIds, err := buffer.UserIDs(b.service)
|
|
if err != nil {
|
|
log.Error(ctx, "Error retrieving userIds from scrobble buffer", "scrobbler", b.service, err)
|
|
return false
|
|
}
|
|
result := true
|
|
for _, userId := range userIds {
|
|
if !b.processUserQueue(ctx, userId) {
|
|
result = false
|
|
}
|
|
}
|
|
return result
|
|
}
|
|
|
|
func (b *bufferedScrobbler) processUserQueue(ctx context.Context, userId string) bool {
|
|
// Scrobbles are drained on a background context that no longer carries the
|
|
// request's authenticated user. Restore it from the buffered userId so that
|
|
// scrobblers relying on the user in the context (e.g. plugins) still get it.
|
|
if user, err := b.ds.User(ctx).Get(userId); err != nil {
|
|
log.Warn(ctx, "Could not load user for buffered scrobble", "userId", userId, "scrobbler", b.service, err)
|
|
} else {
|
|
ctx = request.WithUser(ctx, *user)
|
|
}
|
|
buffer := b.ds.ScrobbleBuffer(ctx)
|
|
for {
|
|
entry, err := buffer.Next(b.service, userId)
|
|
if err != nil {
|
|
log.Error(ctx, "Error reading from scrobble buffer", "scrobbler", b.service, err)
|
|
return false
|
|
}
|
|
if entry == nil {
|
|
return true
|
|
}
|
|
s, ok := b.loader()
|
|
if !ok {
|
|
log.Warn(ctx, "Scrobbler not available, will retry later", "scrobbler", b.service)
|
|
return false
|
|
}
|
|
log.Debug(ctx, "Sending scrobble", "scrobbler", b.service, "track", entry.Title, "artist", entry.Artist)
|
|
err = s.Scrobble(ctx, entry.UserID, Scrobble{
|
|
MediaFile: entry.MediaFile,
|
|
TimeStamp: entry.PlayTime,
|
|
})
|
|
if errors.Is(err, ErrRetryLater) {
|
|
log.Warn(ctx, "Could not send scrobble. Will be retried", "userId", entry.UserID,
|
|
"track", entry.Title, "artist", entry.Artist, "scrobbler", b.service, err)
|
|
return false
|
|
}
|
|
if err != nil {
|
|
log.Error(ctx, "Error sending scrobble to service. Discarding", "scrobbler", b.service,
|
|
"userId", entry.UserID, "artist", entry.Artist, "track", entry.Title, err)
|
|
}
|
|
err = buffer.Dequeue(entry)
|
|
if err != nil {
|
|
log.Error(ctx, "Error removing entry from scrobble buffer", "userId", entry.UserID,
|
|
"track", entry.Title, "artist", entry.Artist, "scrobbler", b.service, err)
|
|
return false
|
|
}
|
|
}
|
|
}
|
|
|
|
var _ Scrobbler = (*bufferedScrobbler)(nil)
|