package library import ( "context" "fmt" "log/slog" "sync" "time" "github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgxpool" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" ) // Loudness backfill (M464 #4995). // // Every track is measured here, the new ones included: the scan only deletes a // changed file's measurement (see scanFile), and this worker measures it again. // Measuring inline in the scan was the first plan and was dropped, because the // analysis decodes the whole file. Added to the scan, a large import would run // several times longer and could pass StuckScanThreshold (1h), at which point // the run is reaped and a second scan started beside it. The fingerprint // backfill is a worker of its own for the same reason. // // Until a track is measured it plays with no gain adjustment, which is how // every track played before normalization existed. // loudnessBackfillTick is how often the worker looks for work. Shorter than // the fingerprint backfill's hour, because a new track plays unleveled until it // is measured; once the library has caught up, a tick is one indexed query. const loudnessBackfillTick = 10 * time.Minute // loudnessBackfillBatch is how many tracks one query hands the worker. const loudnessBackfillBatch = 50 // loudnessBackfillConcurrency is the shipped value of the concurrency setting. // Low for the reason fingerprinting's is: each analysis is a full decode, // competing with playback transcoding and streaming. const loudnessBackfillConcurrency = 2 // BackfillLoudnessResult tallies one pass. type BackfillLoudnessResult struct { Processed int Measured int Silent int // read fine, no block above the gate (settled) Unreadable int // ffmpeg could not decode the file (settled) Inconclusive int // nothing stored; tried again on a later pass } func (r *BackfillLoudnessResult) add(o loudnessOutcome) { r.Processed++ switch o { case loudnessMeasured: r.Measured++ case loudnessSilent: r.Silent++ case loudnessUnreadable: r.Unreadable++ default: r.Inconclusive++ } } // LoudnessBackfillWorker measures every track that has no current measurement. type LoudnessBackfillWorker struct { pool *pgxpool.Pool logger *slog.Logger settings *LoudnessSettingsService tick time.Duration batch int32 albumBatch int32 // analyze is a field so an integration test pins which tracks a pass // touches, not what ffmpeg prints. analyze func(ctx context.Context, path string, durationMs int32) loudnessResult } // NewLoudnessBackfillWorker builds a worker with the production cadence. // settings is shared with the admin API; nil runs on defaults. func NewLoudnessBackfillWorker( pool *pgxpool.Pool, logger *slog.Logger, settings *LoudnessSettingsService, ) *LoudnessBackfillWorker { return &LoudnessBackfillWorker{ pool: pool, logger: logger, settings: settings, tick: loudnessBackfillTick, batch: loudnessBackfillBatch, albumBatch: albumLoudnessBatch, analyze: computeLoudness, } } // Run blocks until ctx is cancelled: one pass at start, then one per tick. func (w *LoudnessBackfillWorker) Run(ctx context.Context) { w.runOnce(ctx) t := time.NewTicker(w.tick) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: w.runOnce(ctx) } } } // runOnce contains a pass so that nothing it does (an error, a panic) can stop // the next tick from firing (rule 157). func (w *LoudnessBackfillWorker) runOnce(ctx context.Context) { defer func() { if r := recover(); r != nil { w.logger.Error("loudness backfill: pass panicked", "panic", r) } }() res, err := w.pass(ctx) if err != nil && ctx.Err() == nil { w.logger.Warn("loudness backfill: pass failed", "err", err, "processed", res.Processed) } if res.Processed > 0 { w.logger.Info("loudness backfill: pass complete", "processed", res.Processed, "measured", res.Measured, "silent", res.Silent, "unreadable", res.Unreadable, "inconclusive", res.Inconclusive) } // Album loudness follows the tracks (#4996). It runs even with analysis // switched off: it decodes nothing, and membership still changes as the // library does. albums, err := w.albumPass(ctx) if err != nil && ctx.Err() == nil { w.logger.Warn("album loudness: pass failed", "err", err, "recomputed", albums.Recomputed) } if albums.Recomputed > 0 || albums.Orphans > 0 { w.logger.Info("album loudness: pass complete", "recomputed", albums.Recomputed, "leveled", albums.Leveled, "waiting", albums.Waiting, "failed", albums.Failed, "orphans", albums.Orphans) } } // pass walks every track needing a measurement once, keyset-paged on id. The // cursor is what lets a pass end: an inconclusive attempt writes no row, so a // file that keeps timing out would otherwise be listed again immediately. // Settings are read before every batch, so switching analysis off ends the // pass and a new concurrency applies to the next batch. func (w *LoudnessBackfillWorker) pass(ctx context.Context) (BackfillLoudnessResult, error) { q := dbq.New(w.pool) var ( res BackfillLoudnessResult mu sync.Mutex ) // The all-zero uuid sorts before every real id. Valid must be true: a NULL // cursor would make "id > NULL" match nothing and every pass a no-op. after := pgtype.UUID{Valid: true} for { if err := ctx.Err(); err != nil { return res, err } cfg := w.settings.Get() if !cfg.Enabled { return res, nil } rows, err := q.ListTracksNeedingLoudness(ctx, dbq.ListTracksNeedingLoudnessParams{ CurrentVersion: loudnessVersion, AfterID: after, BatchLimit: w.batch, }) if err != nil { return res, fmt.Errorf("list tracks needing loudness: %w", err) } if len(rows) == 0 { return res, nil } sem := make(chan struct{}, max(1, int(cfg.BackfillConcurrency))) var wg sync.WaitGroup for _, row := range rows { if ctx.Err() != nil { break } sem <- struct{}{} wg.Add(1) go func(row dbq.ListTracksNeedingLoudnessRow) { defer wg.Done() defer func() { <-sem }() defer func() { if r := recover(); r != nil { w.logger.Error("loudness backfill: track panicked", "path", row.FilePath, "panic", r) } }() outcome := storeLoudness(ctx, q, w.logger, row.ID, row.FilePath, w.analyzeFile(ctx, row.FilePath, row.DurationMs)) mu.Lock() res.add(outcome) mu.Unlock() }(row) } wg.Wait() after = rows[len(rows)-1].ID } } func (w *LoudnessBackfillWorker) analyzeFile(ctx context.Context, path string, durationMs int32) loudnessResult { if w.analyze == nil { return computeLoudness(ctx, path, durationMs) } return w.analyze(ctx, path, durationMs) } // LoudnessCoverage reports how much of the library carries a current // measurement, for the admin gauge. It lives beside the backfill so the // version it counts against is the one the backfill writes. func LoudnessCoverage(ctx context.Context, pool *pgxpool.Pool) (dbq.GetLoudnessCoverageRow, error) { return dbq.New(pool).GetLoudnessCoverage(ctx, loudnessVersion) }