package library import ( "context" "errors" "fmt" "io" "log/slog" "path/filepath" "sync" "testing" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync" ) // TestLoudnessBackfill_Integration pins which tracks a pass measures, what it // stores for each kind of result, that a pass ends, and that the gauge counts // what the passes wrote. func TestLoudnessBackfill_Integration(t *testing.T) { pool := newPool(t) ctx := context.Background() q := dbq.New(pool) dir := t.TempDir() _, album, artist := seedTrack(t, pool, filepath.Join(dir, "unmeasured.mp3")) addTrack := func(name string) dbq.Track { t.Helper() tr, err := q.UpsertTrack(ctx, dbq.UpsertTrackParams{ Title: name, AlbumID: album.ID, ArtistID: artist.ID, DurationMs: 180000, FilePath: filepath.Join(dir, name+".mp3"), FileSize: 100, FileFormat: "mp3", }) if err != nil { t.Fatalf("track %s: %v", name, err) } return tr } lufs := func(v float32) *float32 { return &v } current := addTrack("current") stale := addTrack("stale") missing := addTrack("missing") for _, seed := range []struct { track dbq.Track version int16 }{ {current, loudnessVersion}, {stale, loudnessVersion - 1}, } { if err := q.UpsertTrackLoudness(ctx, dbq.UpsertTrackLoudnessParams{ TrackID: seed.track.ID, IntegratedLufs: lufs(-9), AnalysisVersion: seed.version, }); err != nil { t.Fatalf("seed loudness: %v", err) } } if _, err := pool.Exec(ctx, "UPDATE tracks SET missing_since = now() WHERE id = $1", missing.ID); err != nil { t.Fatalf("mark missing: %v", err) } settings, err := NewLoudnessSettingsService(ctx, pool) if err != nil { t.Fatalf("loudness settings: %v", err) } w := NewLoudnessBackfillWorker(pool, slog.New(slog.NewTextHandler(io.Discard, nil)), settings) // A batch of one forces the keyset cursor across several queries in a pass. w.batch = 1 var mu sync.Mutex calls := map[string]int{} durations := map[string]int32{} start := int16(500) w.analyze = func(_ context.Context, path string, durationMs int32) loudnessResult { name := filepath.Base(path) mu.Lock() calls[name]++ durations[name] = durationMs mu.Unlock() switch name { case "stall.mp3": return loudnessResult{err: fmt.Errorf("ffmpeg: %w", errLoudnessTimeout)} case "corrupt.mp3": return loudnessResult{err: errors.New("ffmpeg exited 1: invalid data")} case "silent.mp3": return loudnessResult{} default: return loudnessResult{ integratedLUFS: lufs(-12.5), truePeakDBTP: lufs(-0.4), rangeLU: lufs(6), hist: blockHistogram{start: start, counts: []int32{3, 0, 7}}, } } } callCount := func(name string) int { mu.Lock() defer mu.Unlock() return calls[name] } // 1. Only the unmeasured track and the stale one are measured, never the // current one or the missing one, and the track's length reaches the // analyzer (it sets the deadline). res, err := w.pass(ctx) if err != nil { t.Fatalf("first pass: %v", err) } if res.Processed != 2 || res.Measured != 2 { t.Fatalf("first pass = %+v, want 2 processed, 2 measured", res) } for name, want := range map[string]int{ "unmeasured.mp3": 1, "stale.mp3": 1, "current.mp3": 0, "missing.mp3": 0, } { if got := callCount(name); got != want { t.Errorf("%s measured %d times, want %d", name, got, want) } } if durations["stale.mp3"] != 180000 { t.Errorf("analyzer got duration %d for stale.mp3, want 180000", durations["stale.mp3"]) } var ( gotLUFS, gotPeak *float32 gotStart *int16 gotHist []int32 gotVersion int16 ) if err := pool.QueryRow(ctx, `SELECT integrated_lufs, true_peak_dbtp, block_hist_start, block_hist, analysis_version FROM track_loudness WHERE track_id = $1`, stale.ID). Scan(&gotLUFS, &gotPeak, &gotStart, &gotHist, &gotVersion); err != nil { t.Fatalf("read stale row back: %v", err) } if gotLUFS == nil || *gotLUFS != -12.5 || gotPeak == nil || *gotPeak != -0.4 || gotStart == nil || *gotStart != start || len(gotHist) != 3 || gotHist[2] != 7 || gotVersion != loudnessVersion { t.Errorf("stale row after re-measuring = lufs %v peak %v hist %v@%v version %d", gotLUFS, gotPeak, gotHist, gotStart, gotVersion) } // Each stored measurement told the sync feed, so caching clients re-read // the track and pick up its gain (#4997). var logged int if err := pool.QueryRow(ctx, `SELECT count(*) FROM library_changes WHERE entity_type = 'track' AND op = 'upsert' AND entity_id = $1`, syncpkg.FormatUUID(stale.ID)).Scan(&logged); err != nil { t.Fatalf("count sync changes: %v", err) } if logged != 1 { t.Errorf("re-measuring stale.mp3 logged %d track changes, want 1", logged) } // 2. A pass after a complete one is a no-op. res, err = w.pass(ctx) if err != nil { t.Fatalf("second pass: %v", err) } if res.Processed != 0 { t.Fatalf("second pass processed %d tracks, want 0", res.Processed) } // 3. A stall is tried once and the pass ends; silence and a corrupt file are // verdicts, stored and not tried again. addTrack("stall") addTrack("corrupt") silent := addTrack("silent") res, err = w.pass(ctx) if err != nil { t.Fatalf("third pass: %v", err) } if res.Processed != 3 || res.Inconclusive != 1 || res.Unreadable != 1 || res.Silent != 1 { t.Fatalf("third pass = %+v, want 3 processed: 1 inconclusive, 1 unreadable, 1 silent", res) } if got := callCount("stall.mp3"); got != 1 { t.Fatalf("stalling file tried %d times in one pass, want exactly 1", got) } var silentHist []int32 var silentUnreadable bool if err := pool.QueryRow(ctx, "SELECT block_hist, unreadable FROM track_loudness WHERE track_id = $1", silent.ID).Scan(&silentHist, &silentUnreadable); err != nil { t.Fatalf("read silent row: %v", err) } if silentHist != nil || silentUnreadable { t.Errorf("silent row = hist %v unreadable %v, want no histogram and readable", silentHist, silentUnreadable) } res, err = w.pass(ctx) if err != nil { t.Fatalf("fourth pass: %v", err) } if res.Processed != 1 || callCount("stall.mp3") != 2 || callCount("corrupt.mp3") != 1 { t.Fatalf("fourth pass = %+v; want only the stalled file retried", res) } // 4. The gauge counts what the passes wrote, and its buckets add up. Six // present tracks: unmeasured, current, stale, stall, corrupt, silent. cov, err := LoudnessCoverage(ctx, pool) if err != nil { t.Fatalf("coverage: %v", err) } if cov.Total != 6 || cov.Measured != 3 || cov.Silent != 1 || cov.Unreadable != 1 || cov.Pending != 1 { t.Errorf("coverage = %+v, want total 6, measured 3, silent 1, unreadable 1, pending 1", cov) } if cov.Measured+cov.Silent+cov.Unreadable+cov.Pending != cov.Total { t.Errorf("coverage buckets %+v do not sum to the total", cov) } // 5. A changed file loses its measurement, so it is measured again. if err := q.DeleteTrackLoudness(ctx, current.ID); err != nil { t.Fatalf("delete loudness: %v", err) } // 6. Switched off, the backfill does nothing, even with work waiting. off := DefaultLoudnessSettings off.Enabled = false if _, err := settings.Set(ctx, off); err != nil { t.Fatalf("switch analysis off: %v", err) } res, err = w.pass(ctx) if err != nil { t.Fatalf("pass with analysis off: %v", err) } if res.Processed != 0 || callCount("current.mp3") != 0 { t.Fatalf("pass with analysis off = %+v, want nothing done", res) } if _, err := settings.Set(ctx, DefaultLoudnessSettings); err != nil { t.Fatalf("switch analysis on: %v", err) } res, err = w.pass(ctx) if err != nil { t.Fatalf("pass after switching back on: %v", err) } if callCount("current.mp3") != 1 { t.Fatalf("changed file measured %d times after switching back on (pass %+v), want 1", callCount("current.mp3"), res) } } // The Go defaults must match the migration's, or a database that cannot be // read would analyze differently from a fresh install. func TestLoudnessSettings_DefaultsMatchMigration(t *testing.T) { pool := newPool(t) s, err := NewLoudnessSettingsService(context.Background(), pool) if err != nil { t.Fatalf("load: %v", err) } got := s.Get() got.UpdatedAt = DefaultLoudnessSettings.UpdatedAt if got != DefaultLoudnessSettings { t.Errorf("migration defaults = %+v, Go defaults = %+v", got, DefaultLoudnessSettings) } }