diff --git a/go.mod b/go.mod index b78e9bb6..2c7393f7 100644 --- a/go.mod +++ b/go.mod @@ -13,6 +13,7 @@ require ( github.com/jackc/pgx/v5 v5.9.2 github.com/stretchr/testify v1.11.1 golang.org/x/crypto v0.51.0 + golang.org/x/sync v0.21.0 gopkg.in/yaml.v3 v3.0.1 ) @@ -25,7 +26,6 @@ require ( github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/robfig/cron/v3 v3.0.1 // indirect github.com/rogpeppe/go-internal v1.14.1 // indirect - golang.org/x/sync v0.21.0 // indirect golang.org/x/sys v0.44.0 // indirect golang.org/x/text v0.39.0 // indirect ) diff --git a/internal/api/admin_loudness.go b/internal/api/admin_loudness.go index 13cf850d..3df88ffc 100644 --- a/internal/api/admin_loudness.go +++ b/internal/api/admin_loudness.go @@ -46,10 +46,15 @@ func (h *handlers) handleGetLoudnessCoverage(w http.ResponseWriter, r *http.Requ type loudnessSettingsBody struct { Enabled bool `json:"enabled"` BackfillConcurrency int32 `json:"backfill_concurrency"` + LeveledCacheMB int32 `json:"leveled_cache_mb"` } func loudnessSettingsBodyOf(s library.LoudnessSettings) loudnessSettingsBody { - return loudnessSettingsBody{Enabled: s.Enabled, BackfillConcurrency: s.BackfillConcurrency} + return loudnessSettingsBody{ + Enabled: s.Enabled, + BackfillConcurrency: s.BackfillConcurrency, + LeveledCacheMB: s.LeveledCacheMB, + } } // handleGetLoudnessSettings implements GET /api/admin/library/loudness-settings. @@ -69,6 +74,7 @@ func (h *handlers) handleUpdateLoudnessSettings(w http.ResponseWriter, r *http.R saved, err := h.loudnessSettings.Set(r.Context(), library.LoudnessSettings{ Enabled: req.Enabled, BackfillConcurrency: req.BackfillConcurrency, + LeveledCacheMB: req.LeveledCacheMB, }) if err != nil { if errors.Is(err, library.ErrLoudnessSettingOutOfRange) { diff --git a/internal/api/api.go b/internal/api/api.go index 03e539bc..55318155 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -7,6 +7,7 @@ package api import ( "log/slog" "math/rand" + "path/filepath" "time" "github.com/go-chi/chi/v5" @@ -66,6 +67,7 @@ func Mount(r chi.Router, pool *pgxpool.Pool, logger *slog.Logger, events *playev reacqSettings: reacqSettings, fingerprintSettings: fpSettings, loudnessSettings: loudSettings, + leveled: newLeveledRenderer(dataDir, loudSettings, logger), librarySize: recommendation.NewLibrarySize(nil), loginGuard: auth.NewLoginGuard(), setupToken: setupToken, @@ -96,6 +98,9 @@ func Mount(r chi.Router, pool *pgxpool.Pool, logger *slog.Logger, events *playev // audio format from the path. The {ext} param is consumed by chi // and ignored by the handler (which keys off {id}). See task #610. api.With(auth.OptionalUser(pool, logger)).Get("/tracks/{id}/stream.{ext}", h.handleGetStream) + // The leveled stream for Sonos/UPnP (M464 #5001): session or a + // leveled token, like the plain stream. + api.With(auth.OptionalUser(pool, logger)).Get("/tracks/{id}/leveled.flac", h.handleGetLeveledStream) api.Group(func(authed chi.Router) { authed.Use(auth.RequireUser(pool, netSettings.Hops)) @@ -343,6 +348,10 @@ type handlers struct { // loudnessSettings is the loudness analysis policy (M464 #4995), the same // instance the loudness backfill reads. Nil serves the defaults. loudnessSettings *library.LoudnessSettingsService + // leveled renders the leveled streams handed to Sonos/UPnP speakers + // (M464 #5001). Nil when its cache directory cannot be made: a level + // request then gets the plain stream. + leveled *library.LeveledRenderer // setupToken must accompany the first registration while no users exist // (see auth.SetupToken). requireSetupToken is set by Mount, the only // production constructor; tests that build handlers directly leave it @@ -370,3 +379,14 @@ type handlers struct { // anything a client mints), which is the desired slice-1 default. streamSecret []byte } + +// newLeveledRenderer makes the leveled-stream renderer, caching under the +// data directory. A failure is logged and leaves leveling off for speakers. +func newLeveledRenderer(dataDir string, settings *library.LoudnessSettingsService, logger *slog.Logger) *library.LeveledRenderer { + r, err := library.NewLeveledRenderer(filepath.Join(dataDir, "leveled-cache"), settings, logger) + if err != nil { + logger.Error("api: leveled streams unavailable", "err", err) + return nil + } + return r +} diff --git a/internal/api/cast_token.go b/internal/api/cast_token.go index 0e18e869..a1594874 100644 --- a/internal/api/cast_token.go +++ b/internal/api/cast_token.go @@ -8,6 +8,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/apierror" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/library" ) const ( @@ -19,6 +20,12 @@ const ( type castTokenRequest struct { TrackID string `json:"trackId"` ExpSeconds int `json:"expSeconds,omitempty"` + // Level asks for the leveled stream (M464 #5001): the track rendered at + // the user's loudness gain. AsAlbum says the track is playing as part of + // its album in order, which picks album gain in auto mode; only the + // client holding the queue knows it. + Level bool `json:"level,omitempty"` + AsAlbum bool `json:"asAlbum,omitempty"` } type castTokenResponse struct { @@ -31,6 +38,10 @@ type castTokenResponse struct { // `` and `` without a follow-up round trip. MIME string `json:"mime"` Title string `json:"title"` + // Leveled is true when URL is the leveled stream. A level request still + // gets the plain stream when there is nothing to change: leveling off, + // the track not yet measured, or a gain of 0. + Leveled bool `json:"leveled"` } // mimeForFormat returns the audio MIME type for a cast (Sonos/UPnP) URL. @@ -89,7 +100,8 @@ func extForFormat(format string) string { // Part of the output-picker UPnP slice. See // docs/superpowers/specs/2026-06-03-android-output-picker-upnp-design.md. func (h *handlers) handleCastStreamToken(w http.ResponseWriter, r *http.Request) { - if _, ok := requireUser(w, r); !ok { + user, ok := requireUser(w, r) + if !ok { return } var req castTokenRequest @@ -112,6 +124,27 @@ func (h *handlers) handleCastStreamToken(w http.ResponseWriter, r *http.Request) expSec := clampExpSeconds(req.ExpSeconds) exp := time.Now().Add(time.Duration(expSec) * time.Second).Unix() token := SignStreamToken(h.streamSecret, req.TrackID, exp) + path := streamURLWithExt(trackUUID, extForFormat(track.FileFormat)) + + "?token=" + token + "&exp=" + strconv.FormatInt(exp, 10) + mime := mimeForFormat(track.FileFormat) + leveled := false + if req.Level && h.leveled != nil { + g, err := h.leveledGainFor(r.Context(), user.ID, trackUUID, req.AsAlbum) + if err != nil { + // The plain stream still plays; only the leveling is lost. + h.logger.Warn("cast token: leveled gain lookup failed", "track", req.TrackID, "err", err) + } else if !g.Unity() { + token = SignLeveledStreamToken(h.streamSecret, req.TrackID, exp, g) + path = leveledStreamPath(trackUUID) + leveledQuery(g, token, exp) + mime = "audio/flac" + leveled = true + // The speaker fetches the URL shortly; start rendering now so + // the fetch finds the file ready. + h.leveled.Prerender(library.LeveledSource{ + TrackID: req.TrackID, Path: track.FilePath, DurationMs: track.DurationMs, + }, g) + } + } // Behind a TLS-terminating reverse proxy, r.TLS is nil even though // the public-facing URL is https://. UPnP devices (Sonos especially) @@ -131,18 +164,16 @@ func (h *handlers) handleCastStreamToken(w http.ResponseWriter, r *http.Request) if h := r.Header.Get("X-Forwarded-Host"); h != "" { host = h } - // Include the file extension in the path so Sonos's URL probe sees a + // The path carries a file extension so Sonos's URL probe sees a // recognizable audio file. Without it, Sonos reports TrackDuration=0 // and seeks past 0s land "after the end" -> early track-skip. - url := scheme + "://" + host + streamURLWithExt(trackUUID, extForFormat(track.FileFormat)) + - "?token=" + token + "&exp=" + strconv.FormatInt(exp, 10) - writeJSON(w, http.StatusOK, castTokenResponse{ - Token: token, - Exp: exp, - URL: url, - MIME: mimeForFormat(track.FileFormat), - Title: track.Title, + Token: token, + Exp: exp, + URL: scheme + "://" + host + path, + MIME: mime, + Title: track.Title, + Leveled: leveled, }) } diff --git a/internal/api/leveled_stream.go b/internal/api/leveled_stream.go new file mode 100644 index 00000000..6c01f708 --- /dev/null +++ b/internal/api/leveled_stream.go @@ -0,0 +1,168 @@ +package api + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "net/http" + "os" + "strconv" + "time" + + "github.com/go-chi/chi/v5" + "github.com/jackc/pgx/v5/pgtype" + + "git.fabledsword.com/bvandeusen/minstrel/internal/apierror" + "git.fabledsword.com/bvandeusen/minstrel/internal/auth" + "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/library" +) + +// The leveled stream (M464 #5001): a track rendered with the user's loudness +// gain applied, for the Sonos and UPnP speakers that fetch their own audio. +// See internal/library/leveled.go for how it is rendered. + +// SignLeveledStreamToken signs a leveled stream URL. The gain is part of what +// is signed, so a speaker's URL cannot be edited into a different render. +// The message cannot collide with SignStreamToken's "|": a plain +// token does not open a leveled stream, nor the reverse. +func SignLeveledStreamToken(secret []byte, trackID string, exp int64, g library.LeveledGain) string { + mac := hmac.New(sha256.New, secret) + lim := 0 + if g.Limiter { + lim = 1 + } + _, _ = fmt.Fprintf(mac, "%s|%d|leveled|%d|%d", trackID, exp, g.CentiDB, lim) + return hex.EncodeToString(mac.Sum(nil)) +} + +// VerifyLeveledStreamToken checks a token from SignLeveledStreamToken, and +// that it has not expired. +func VerifyLeveledStreamToken(secret []byte, trackID string, exp int64, g library.LeveledGain, token string) bool { + if time.Now().Unix() > exp { + return false + } + return hmac.Equal([]byte(SignLeveledStreamToken(secret, trackID, exp, g)), []byte(token)) +} + +// leveledStreamPath is the leveled stream's path. It ends in .flac because +// Sonos reads the format from the URL's extension (task #610). +func leveledStreamPath(trackID pgtype.UUID) string { + return "/api/tracks/" + uuidToString(trackID) + "/leveled.flac" +} + +// leveledQuery is the query string a leveled URL carries. +func leveledQuery(g library.LeveledGain, token string, exp int64) string { + lim := "0" + if g.Limiter { + lim = "1" + } + return "?g=" + strconv.Itoa(g.CentiDB) + "&lim=" + lim + + "&token=" + token + "&exp=" + strconv.FormatInt(exp, 10) +} + +// leveledGainFor is the render request for one track under the user's +// preference. Unity when leveling is off or the track is unmeasured, which +// the caller answers with the plain stream. +func (h *handlers) leveledGainFor(ctx context.Context, userID, trackID pgtype.UUID, asAlbum bool) (library.LeveledGain, error) { + q := dbq.New(h.pool) + prefs, err := library.LoadNormalizationPrefs(ctx, q, userID) + if err != nil { + return library.LeveledGain{}, err + } + gains, err := library.ReplayGainForTracks(ctx, q, []pgtype.UUID{trackID}) + if err != nil { + return library.LeveledGain{}, err + } + return library.NewLeveledGain(library.LeveledGainDB(prefs, gains[trackID], asAlbum), prefs.Boost), nil +} + +// parseLeveledGain reads ?g= and ?lim= from a leveled URL. +func parseLeveledGain(r *http.Request) (library.LeveledGain, bool) { + c, err := strconv.Atoi(r.URL.Query().Get("g")) + if err != nil { + return library.LeveledGain{}, false + } + lim := r.URL.Query().Get("lim") + if lim != "0" && lim != "1" { + return library.LeveledGain{}, false + } + g := library.LeveledGain{CentiDB: c, Limiter: lim == "1"} + return g, g.Valid() +} + +// handleGetLeveledStream implements GET /api/tracks/{id}/leveled.flac. It +// accepts a session, like the plain stream, or a leveled token: the gain in +// the query must be the one the token was signed for. +func (h *handlers) handleGetLeveledStream(w http.ResponseWriter, r *http.Request) { + rawID := chi.URLParam(r, "id") + g, ok := parseLeveledGain(r) + if !h.leveledAuthOk(r, rawID, g, ok) { + writeErr(w, apierror.ErrUnauthorized) + return + } + if !ok { + writeErr(w, apierror.BadRequest("invalid_gain", "g and lim must describe a valid gain")) + return + } + if h.leveled == nil { + writeErr(w, &apierror.Error{Status: http.StatusServiceUnavailable, Code: "leveling_unavailable", + Message: "leveled streams are not available on this server"}) + return + } + track, apiErr := resolveByID(r, "id", dbq.New(h.pool).GetTrackByID, "track") + if apiErr != nil { + writeErr(w, apiErr) + return + } + path, err := h.leveled.Path(r.Context(), library.LeveledSource{ + TrackID: rawID, Path: track.FilePath, DurationMs: track.DurationMs, + }, g) + switch { + case errors.Is(err, library.ErrLeveledSourceMissing): + writeErr(w, &apierror.Error{Status: http.StatusNotFound, Code: "not_found", Message: "track file not found"}) + return + case r.Context().Err() != nil: + return // the speaker hung up; the render carries on for its retry + case err != nil: + writeErrWithLog(w, h.logger, "leveled stream: render failed", apierror.InternalMsg("render failed", err)) + return + } + f, err := os.Open(path) + if err != nil { + // Evicted between render and open: rare, and the speaker retries. + writeErrWithLog(w, h.logger, "leveled stream: open render", apierror.InternalMsg("server error", err)) + return + } + defer func() { _ = f.Close() }() + info, err := f.Stat() + if err != nil { + writeErrWithLog(w, h.logger, "leveled stream: stat render", apierror.InternalMsg("server error", err)) + return + } + w.Header().Set("Content-Type", "audio/flac") + w.Header().Set("Accept-Ranges", "bytes") + w.Header().Set("Cache-Control", "private, max-age=86400") + http.ServeContent(w, r, "leveled.flac", info.ModTime(), f) +} + +// leveledAuthOk mirrors streamAuthOk: a session, or a token signed over this +// track and this gain. A malformed gain fails the token path, since there is +// nothing it could have been signed over. +func (h *handlers) leveledAuthOk(r *http.Request, trackID string, g library.LeveledGain, gainOK bool) bool { + if _, ok := auth.UserFromContext(r.Context()); ok { + return true + } + if !gainOK { + return false + } + tok := r.URL.Query().Get("token") + exp, err := strconv.ParseInt(r.URL.Query().Get("exp"), 10, 64) + if tok == "" || err != nil { + return false + } + return VerifyLeveledStreamToken(h.streamSecret, trackID, exp, g, tok) +} diff --git a/internal/api/leveled_stream_test.go b/internal/api/leveled_stream_test.go new file mode 100644 index 00000000..9bab2feb --- /dev/null +++ b/internal/api/leveled_stream_test.go @@ -0,0 +1,161 @@ +package api + +import ( + "bytes" + "context" + "encoding/json" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "strconv" + "strings" + "testing" + "time" + + "github.com/go-chi/chi/v5" + + "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/library" +) + +func TestLeveledStreamToken(t *testing.T) { + secret := []byte("leveled-secret") + exp := time.Now().Unix() + 3600 + g := library.LeveledGain{CentiDB: -550} + tok := SignLeveledStreamToken(secret, "track-1", exp, g) + + if !VerifyLeveledStreamToken(secret, "track-1", exp, g, tok) { + t.Fatal("round trip failed") + } + // The gain is part of what is signed: a speaker URL edited to another + // gain, or to switch the limiter on, no longer verifies. + if VerifyLeveledStreamToken(secret, "track-1", exp, library.LeveledGain{CentiDB: 1200}, tok) { + t.Error("token verified for a different gain") + } + if VerifyLeveledStreamToken(secret, "track-1", exp, library.LeveledGain{CentiDB: -550, Limiter: true}, tok) { + t.Error("token verified with the limiter switched on") + } + if VerifyLeveledStreamToken(secret, "track-2", exp, g, tok) { + t.Error("token verified for another track") + } + if VerifyLeveledStreamToken(secret, "track-1", time.Now().Unix()-1, g, SignLeveledStreamToken(secret, "track-1", time.Now().Unix()-1, g)) { + t.Error("expired token verified") + } + // Plain and leveled tokens are not interchangeable. + plain := SignStreamToken(secret, "track-1", exp) + if VerifyLeveledStreamToken(secret, "track-1", exp, library.LeveledGain{}, plain) { + t.Error("a plain stream token opened a leveled stream") + } + if VerifyStreamToken(secret, "track-1", exp, tok) { + t.Error("a leveled token opened the plain stream") + } +} + +func TestParseLeveledGain(t *testing.T) { + for _, c := range []struct { + query string + want library.LeveledGain + ok bool + }{ + {"g=-550&lim=0", library.LeveledGain{CentiDB: -550}, true}, + {"g=320&lim=1", library.LeveledGain{CentiDB: 320, Limiter: true}, true}, + {"g=1201&lim=0", library.LeveledGain{}, false}, + {"g=abc&lim=0", library.LeveledGain{}, false}, + {"g=100&lim=yes", library.LeveledGain{}, false}, + {"lim=0", library.LeveledGain{}, false}, + } { + got, ok := parseLeveledGain(httptest.NewRequest(http.MethodGet, "/x?"+c.query, nil)) + if ok != c.ok || (ok && got != c.want) { + t.Errorf("%s: got %+v ok=%v, want %+v ok=%v", c.query, got, ok, c.want, c.ok) + } + } +} + +func TestCastStreamToken_Leveled(t *testing.T) { + h, pool := testHandlers(t) + h.streamSecret = []byte("leveled-cast-secret") + r, err := library.NewLeveledRenderer(t.TempDir(), nil, slog.New(slog.NewTextHandler(io.Discard, nil))) + if err != nil { + t.Fatal(err) + } + h.leveled = r + user := seedUser(t, pool, "lev", "pw", false) + track, _ := seedTrackForRemoveTest(t, h, "lev", "", "") + trackID := uuidToString(track.ID) + lufs, peak := float32(-12.5), float32(-1) + if err := dbq.New(pool).UpsertTrackLoudness(context.Background(), dbq.UpsertTrackLoudnessParams{ + TrackID: track.ID, IntegratedLufs: &lufs, TruePeakDbtp: &peak, AnalysisVersion: 1, + }); err != nil { + t.Fatalf("seed loudness: %v", err) + } + mint := func(req castTokenRequest) castTokenResponse { + t.Helper() + body, _ := json.Marshal(req) + w := httptest.NewRecorder() + h.handleCastStreamToken(w, withUser(httptest.NewRequest(http.MethodPost, "/api/cast/stream-token", bytes.NewReader(body)), user)) + if w.Code != http.StatusOK { + t.Fatalf("status %d: %s", w.Code, w.Body.String()) + } + var resp castTokenResponse + if err := json.NewDecoder(w.Body).Decode(&resp); err != nil { + t.Fatal(err) + } + return resp + } + + // Default preference (auto, -18): the track measures -12.5 LUFS, so it is + // cut by 5.5 dB, and the URL and its token say exactly that. + resp := mint(castTokenRequest{TrackID: trackID, Level: true}) + if !resp.Leveled || resp.MIME != "audio/flac" || + !strings.Contains(resp.URL, "/api/tracks/"+trackID+"/leveled.flac?g=-550&lim=0&token=") { + t.Fatalf("leveled mint = %+v", resp) + } + if !VerifyLeveledStreamToken(h.streamSecret, trackID, resp.Exp, library.LeveledGain{CentiDB: -550}, resp.Token) { + t.Fatal("leveled token does not verify for the gain in the URL") + } + + // Not asked for: the plain stream, as before. + if resp := mint(castTokenRequest{TrackID: trackID}); resp.Leveled || strings.Contains(resp.URL, "leveled") { + t.Fatalf("unleveled mint = %+v", resp) + } + + // Leveling switched off: nothing to render, so the plain stream. + if _, err := library.SaveNormalizationPrefs(context.Background(), dbq.New(pool), user.ID, + library.NormalizationPrefs{Mode: "off", TargetLUFS: -18, Boost: "headroom"}); err != nil { + t.Fatal(err) + } + if resp := mint(castTokenRequest{TrackID: trackID, Level: true}); resp.Leveled { + t.Fatalf("mint with leveling off = %+v, want the plain stream", resp) + } +} + +func TestGetLeveledStream_RefusesUnsignedGains(t *testing.T) { + h, _ := testHandlers(t) + h.streamSecret = []byte("leveled-get-secret") + router := chi.NewRouter() + router.Get("/api/tracks/{id}/leveled.flac", h.handleGetLeveledStream) + id := nonExistentTrackUUID + exp := time.Now().Unix() + 3600 + tok := SignLeveledStreamToken(h.streamSecret, id, exp, library.LeveledGain{CentiDB: -300}) + get := func(query string) int { + w := httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest(http.MethodGet, "/api/tracks/"+id+"/leveled.flac?"+query, nil)) + return w.Code + } + e := strconv.FormatInt(exp, 10) + for name, q := range map[string]string{ + "no token": "g=-300&lim=0", + "edited gain": "g=1200&lim=0&token=" + tok + "&exp=" + e, + "limiter on": "g=-300&lim=1&token=" + tok + "&exp=" + e, + "invalid gain": "g=99999&lim=0&token=" + tok + "&exp=" + e, + } { + if code := get(q); code != http.StatusUnauthorized { + t.Errorf("%s: status %d, want 401", name, code) + } + } + // The signed gain gets past auth to the track lookup. + if code := get("g=-300&lim=0&token=" + tok + "&exp=" + e); code == http.StatusUnauthorized { + t.Errorf("the signed gain was refused") + } +} diff --git a/internal/db/dbq/loudness.sql.go b/internal/db/dbq/loudness.sql.go index 59d3c143..c7d1df1d 100644 --- a/internal/db/dbq/loudness.sql.go +++ b/internal/db/dbq/loudness.sql.go @@ -127,7 +127,7 @@ func (q *Queries) GetLoudnessCoverage(ctx context.Context, currentVersion int16) } const getLoudnessSettings = `-- name: GetLoudnessSettings :one -SELECT id, enabled, backfill_concurrency, updated_at FROM loudness_settings WHERE id = true +SELECT id, enabled, backfill_concurrency, updated_at, leveled_cache_mb FROM loudness_settings WHERE id = true ` func (q *Queries) GetLoudnessSettings(ctx context.Context) (LoudnessSetting, error) { @@ -138,6 +138,7 @@ func (q *Queries) GetLoudnessSettings(ctx context.Context) (LoudnessSetting, err &i.Enabled, &i.BackfillConcurrency, &i.UpdatedAt, + &i.LeveledCacheMb, ) return i, err } @@ -388,26 +389,29 @@ const updateLoudnessSettings = `-- name: UpdateLoudnessSettings :one UPDATE loudness_settings SET enabled = $1, backfill_concurrency = $2, + leveled_cache_mb = $3, updated_at = now() WHERE id = true -RETURNING id, enabled, backfill_concurrency, updated_at +RETURNING id, enabled, backfill_concurrency, updated_at, leveled_cache_mb ` type UpdateLoudnessSettingsParams struct { Enabled bool BackfillConcurrency int32 + LeveledCacheMb int32 } -// Whole-row write from the admin card; migration 0065's CHECK is the backstop +// Whole-row write from the admin card; the migrations' CHECKs are the backstop // behind the service's own validation. func (q *Queries) UpdateLoudnessSettings(ctx context.Context, arg UpdateLoudnessSettingsParams) (LoudnessSetting, error) { - row := q.db.QueryRow(ctx, updateLoudnessSettings, arg.Enabled, arg.BackfillConcurrency) + row := q.db.QueryRow(ctx, updateLoudnessSettings, arg.Enabled, arg.BackfillConcurrency, arg.LeveledCacheMb) var i LoudnessSetting err := row.Scan( &i.ID, &i.Enabled, &i.BackfillConcurrency, &i.UpdatedAt, + &i.LeveledCacheMb, ) return i, err } diff --git a/internal/db/dbq/models.go b/internal/db/dbq/models.go index 90a6fb80..e62b9121 100644 --- a/internal/db/dbq/models.go +++ b/internal/db/dbq/models.go @@ -432,6 +432,7 @@ type LoudnessSetting struct { Enabled bool BackfillConcurrency int32 UpdatedAt pgtype.Timestamptz + LeveledCacheMb int32 } type MissingReacquisition struct { diff --git a/internal/db/migrations/0068_leveled_stream_cache.down.sql b/internal/db/migrations/0068_leveled_stream_cache.down.sql new file mode 100644 index 00000000..59ab4d17 --- /dev/null +++ b/internal/db/migrations/0068_leveled_stream_cache.down.sql @@ -0,0 +1,3 @@ +ALTER TABLE loudness_settings + DROP CONSTRAINT loudness_settings_leveled_cache_range, + DROP COLUMN leveled_cache_mb; diff --git a/internal/db/migrations/0068_leveled_stream_cache.up.sql b/internal/db/migrations/0068_leveled_stream_cache.up.sql new file mode 100644 index 00000000..941bd128 --- /dev/null +++ b/internal/db/migrations/0068_leveled_stream_cache.up.sql @@ -0,0 +1,12 @@ +-- 0068_leveled_stream_cache.up.sql — the size cap of the leveled-stream render +-- cache (Scribe milestone #464, #5001). +-- +-- Sonos and other UPnP renderers fetch a track's URL themselves, so the phone +-- cannot level what they play. The server renders a gain-applied FLAC for them +-- instead and keeps it on disk, oldest-used first out once the cache is over +-- this size. An operator setting (rule 25), on the loudness settings row it +-- belongs with. +ALTER TABLE loudness_settings + ADD COLUMN leveled_cache_mb integer NOT NULL DEFAULT 2048, + ADD CONSTRAINT loudness_settings_leveled_cache_range + CHECK (leveled_cache_mb >= 256 AND leveled_cache_mb <= 65536); diff --git a/internal/db/queries/loudness.sql b/internal/db/queries/loudness.sql index af33fb1b..90a92c2d 100644 --- a/internal/db/queries/loudness.sql +++ b/internal/db/queries/loudness.sql @@ -67,11 +67,12 @@ SELECT count(*)::bigint AS total, SELECT * FROM loudness_settings WHERE id = true; -- name: UpdateLoudnessSettings :one --- Whole-row write from the admin card; migration 0065's CHECK is the backstop +-- Whole-row write from the admin card; the migrations' CHECKs are the backstop -- behind the service's own validation. UPDATE loudness_settings SET enabled = sqlc.arg(enabled), backfill_concurrency = sqlc.arg(backfill_concurrency), + leveled_cache_mb = sqlc.arg(leveled_cache_mb), updated_at = now() WHERE id = true RETURNING *; diff --git a/internal/library/leveled.go b/internal/library/leveled.go new file mode 100644 index 00000000..4fb982b7 --- /dev/null +++ b/internal/library/leveled.go @@ -0,0 +1,339 @@ +package library + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "io/fs" + "log/slog" + "math" + "os" + "os/exec" + "path/filepath" + "slices" + "strconv" + "strings" + "sync" + "time" + + "golang.org/x/sync/singleflight" +) + +// Leveled streams (M464 #5001). +// +// Sonos and other UPnP renderers fetch a track's URL themselves, so the phone +// that sent them there cannot change what they play. For them the server +// renders the track with its gain already applied: a FLAC, written to a cache +// directory and served as a plain file, so Content-Length, Range and seeking +// behave as for the original. Tags are stripped; the renderer is told the +// metadata in the DIDL-Lite it is handed, and an embedded ReplayGain tag would +// invite it to adjust a second time. + +// The same limits as the web and Android players (web/src/lib/player/gain.ts, +// android .../player/gain/GainMath.kt): headroom mode stops a boost 1 dB under +// the true peak, and no boost exceeds 12 dB. +const ( + leveledPeakCeilingDBTP = -1.0 + leveledMaxBoostDB = 12.0 + // The cut is bounded too: a gain is a request parameter, and a value no + // measurement could produce is refused rather than rendered. + leveledMaxCutDB = -60.0 +) + +// LeveledGainDB is the gain, in dB, a track gets under prefs: 0 when leveling +// is off or the track has not been measured. asAlbum says whether the track is +// being played as part of its album in order, which only the client that +// holds the queue can know. +func LeveledGainDB(prefs NormalizationPrefs, g ReplayGain, asAlbum bool) float64 { + if prefs.Mode == "off" { + return 0 + } + wantAlbum := prefs.Mode == "album" || (prefs.Mode == "auto" && asAlbum) + // Album gain falls back to track gain while the album is still being + // measured; track gain never falls back to album gain. + gain, peak := g.TrackGain, g.TrackPeak + if wantAlbum && g.AlbumGain != nil { + gain, peak = g.AlbumGain, g.AlbumPeak + } + if gain == nil { + return 0 + } + db := float64(*gain) + float64(prefs.TargetLUFS) - ReplayGainReferenceLUFS + if prefs.Boost == "headroom" && peak != nil && *peak > 0 { + db = math.Min(db, leveledPeakCeilingDBTP-20*math.Log10(float64(*peak))) + } + return math.Min(db, leveledMaxBoostDB) +} + +// LeveledGain is a render request: a gain in hundredths of a dB, and whether +// a limiter holds the peaks of a boost. Hundredths, so it travels in a URL and +// a signature exactly. +type LeveledGain struct { + CentiDB int + Limiter bool +} + +// NewLeveledGain rounds db to a render request. The limiter is only ever +// engaged for a boost: a cut cannot raise a peak. +func NewLeveledGain(db float64, boost string) LeveledGain { + c := int(math.Round(db * 100)) + return LeveledGain{CentiDB: c, Limiter: boost == "limiter" && c > 0} +} + +// Unity reports whether the request would reproduce the original. +func (g LeveledGain) Unity() bool { return g.CentiDB == 0 && !g.Limiter } + +// Valid reports whether the gain is one a measurement could have produced. +func (g LeveledGain) Valid() bool { + return g.CentiDB >= int(leveledMaxCutDB*100) && g.CentiDB <= int(leveledMaxBoostDB*100) +} + +// leveledFilter is the ffmpeg audio filter for g. The limiter's ceiling is +// -1 dBFS (0.891 linear); level=disabled stops alimiter from raising the +// output back to full scale afterwards, which would undo the leveling. +func leveledFilter(g LeveledGain) string { + f := "volume=" + strconv.FormatFloat(float64(g.CentiDB)/100, 'f', 2, 64) + "dB" + if g.Limiter { + f += ",alimiter=limit=0.891:level=disabled" + } + return f +} + +// leveledSource is what the renderer needs to know about the original file. +type leveledSource struct { + SampleRate int + // Bits is the source's sample depth; 0 for a lossy source, which has none. + Bits int +} + +// leveledRenderArgs is the ffmpeg command line rendering src to dst. Output is +// FLAC at the source's depth (24-bit for a hi-res source, 16 otherwise) and at +// most 48 kHz: renderers that take FLAC take those, and few take more. +func leveledRenderArgs(src, dst string, g LeveledGain, s leveledSource) []string { + args := []string{ + "-hide_banner", "-nostdin", "-nostats", "-loglevel", "error", + "-i", src, + "-map", "0:a:0", + "-map_metadata", "-1", + "-af", leveledFilter(g), + "-c:a", "flac", + } + if s.Bits > 16 { + args = append(args, "-sample_fmt", "s32", "-bits_per_raw_sample", "24") + } else { + args = append(args, "-sample_fmt", "s16") + } + if s.SampleRate > 48000 { + args = append(args, "-ar", "48000") + } + return append(args, "-f", "flac", "-y", dst) +} + +// leveledCacheKey names the rendered file. It changes with the source file's +// size and modification time, so a replaced file is never served from an old +// render, and with leveledRenderVersion, so a change to how renders are made +// retires the old ones. +func leveledCacheKey(trackID string, info fs.FileInfo, g LeveledGain) string { + h := sha256.New() + _, _ = fmt.Fprintf(h, "%d|%s|%d|%d|%d|%t", + leveledRenderVersion, trackID, info.Size(), info.ModTime().UnixNano(), g.CentiDB, g.Limiter) + return hex.EncodeToString(h.Sum(nil))[:32] + ".flac" +} + +const leveledRenderVersion = 1 + +// LeveledSource identifies the original a render is made from. +type LeveledSource struct { + TrackID string + Path string + DurationMs int32 +} + +// LeveledRenderer renders leveled FLACs into a cache directory and keeps that +// directory under the size the loudness settings allow. Safe for concurrent +// use: renders of the same file and gain are coalesced into one. +type LeveledRenderer struct { + dir string + settings *LoudnessSettingsService + logger *slog.Logger + group singleflight.Group + evictMu sync.Mutex + + // render and probe are ffmpeg and ffprobe; tests replace them. + render func(ctx context.Context, src, dst string, g LeveledGain, s leveledSource) error + probe func(ctx context.Context, src string) (leveledSource, error) +} + +// NewLeveledRenderer renders into dir, creating it if needed. +func NewLeveledRenderer(dir string, settings *LoudnessSettingsService, logger *slog.Logger) (*LeveledRenderer, error) { + if err := os.MkdirAll(dir, 0o750); err != nil { + return nil, fmt.Errorf("leveled cache: %w", err) + } + return &LeveledRenderer{ + dir: dir, + settings: settings, + logger: logger, + render: ffmpegRenderLeveled, + probe: ffprobeLeveledSource, + }, nil +} + +// ErrLeveledSourceMissing is returned when the original file cannot be read. +var ErrLeveledSourceMissing = errors.New("leveled: source file missing") + +// Path returns the rendered file for src at gain g, rendering it first if it +// is not cached. A caller arriving while the same render runs waits for it +// rather than starting another. +func (r *LeveledRenderer) Path(ctx context.Context, src LeveledSource, g LeveledGain) (string, error) { + info, err := os.Stat(src.Path) + if err != nil { + return "", fmt.Errorf("%w: %w", ErrLeveledSourceMissing, err) + } + dst := filepath.Join(r.dir, leveledCacheKey(src.TrackID, info, g)) + if _, err := os.Stat(dst); err == nil { + // A hit counts as a use, for eviction. + now := time.Now() + _ = os.Chtimes(dst, now, now) + return dst, nil + } + // The render runs detached from the caller: a renderer that gives up on + // one request (Sonos retries quickly) must not cancel the render a second + // request is about to wait on. + ch := r.group.DoChan(dst, func() (any, error) { + rctx, cancel := context.WithTimeout(context.Background(), loudnessTimeout(src.DurationMs)) + defer cancel() + return dst, r.renderTo(rctx, src.Path, dst, g) + }) + select { + case res := <-ch: + if res.Err != nil { + return "", res.Err + } + return dst, nil + case <-ctx.Done(): + return "", ctx.Err() + } +} + +// Prerender starts rendering src at g in the background, so the renderer's +// fetch finds it ready. Used when a leveled URL is handed out: the speaker +// asks for the next track shortly before it plays. +func (r *LeveledRenderer) Prerender(src LeveledSource, g LeveledGain) { + go func() { + ctx, cancel := context.WithTimeout(context.Background(), loudnessTimeout(src.DurationMs)) + defer cancel() + if _, err := r.Path(ctx, src, g); err != nil { + r.logger.Warn("leveled: prerender failed", "track", src.TrackID, "err", err) + } + }() +} + +func (r *LeveledRenderer) renderTo(ctx context.Context, src, dst string, g LeveledGain) error { + s, err := r.probe(ctx, src) + if err != nil { + return fmt.Errorf("leveled: probe: %w", err) + } + // Rendered beside the destination and renamed into place, so a reader + // never sees a half-written file and a failed render leaves nothing. + tmp := dst + ".part" + if err := r.render(ctx, src, tmp, g, s); err != nil { + _ = os.Remove(tmp) + return fmt.Errorf("leveled: render: %w", err) + } + if err := os.Rename(tmp, dst); err != nil { + _ = os.Remove(tmp) + return fmt.Errorf("leveled: %w", err) + } + r.evict(dst) + return nil +} + +// evict removes the least recently used renders until the cache fits the +// configured size. keep is never removed: it is the render about to be served. +func (r *LeveledRenderer) evict(keep string) { + r.evictMu.Lock() + defer r.evictMu.Unlock() + limit := int64(r.settings.Get().LeveledCacheMB) << 20 + entries, err := os.ReadDir(r.dir) + if err != nil { + r.logger.Warn("leveled: read cache", "err", err) + return + } + type cached struct { + path string + size int64 + used time.Time + } + var files []cached + var total int64 + for _, e := range entries { + if e.IsDir() || !strings.HasSuffix(e.Name(), ".flac") { + continue + } + info, err := e.Info() + if err != nil { + continue + } + files = append(files, cached{filepath.Join(r.dir, e.Name()), info.Size(), info.ModTime()}) + total += info.Size() + } + slices.SortFunc(files, func(a, b cached) int { return a.used.Compare(b.used) }) + for _, f := range files { + if total <= limit { + return + } + if f.path == keep { + continue + } + // A file still being sent stays readable: the open descriptor holds it. + if err := os.Remove(f.path); err == nil { + total -= f.size + } + } +} + +func ffmpegRenderLeveled(ctx context.Context, src, dst string, g LeveledGain, s leveledSource) error { + out, err := exec.CommandContext(ctx, "ffmpeg", leveledRenderArgs(src, dst, g, s)...).CombinedOutput() + if err != nil { + return fmt.Errorf("ffmpeg: %w: %s", err, strings.TrimSpace(string(out))) + } + return nil +} + +// ffprobeLeveledSource reads the first audio stream's rate and depth. +// bits_per_raw_sample is set for lossless codecs and 0 or N/A for lossy ones. +func ffprobeLeveledSource(ctx context.Context, src string) (leveledSource, error) { + out, err := exec.CommandContext(ctx, "ffprobe", + "-v", "error", "-select_streams", "a:0", + "-show_entries", "stream=sample_rate,bits_per_raw_sample,bits_per_sample", + "-of", "default=noprint_wrappers=1", src).Output() + if err != nil { + return leveledSource{}, fmt.Errorf("ffprobe: %w", err) + } + return parseLeveledProbe(string(out)), nil +} + +// parseLeveledProbe reads ffprobe's key=value lines. Unknown or N/A values +// read as 0: a 16-bit, at-most-48 kHz render, which suits every source. +func parseLeveledProbe(out string) leveledSource { + var s leveledSource + for _, line := range strings.Split(out, "\n") { + k, v, ok := strings.Cut(strings.TrimSpace(line), "=") + if !ok { + continue + } + n, err := strconv.Atoi(v) + if err != nil { + continue + } + switch k { + case "sample_rate": + s.SampleRate = n + case "bits_per_raw_sample", "bits_per_sample": + s.Bits = max(s.Bits, n) + } + } + return s +} diff --git a/internal/library/leveled_test.go b/internal/library/leveled_test.go new file mode 100644 index 00000000..ff18a2e1 --- /dev/null +++ b/internal/library/leveled_test.go @@ -0,0 +1,244 @@ +package library + +import ( + "context" + "errors" + "io" + "log/slog" + "math" + "os" + "path/filepath" + "slices" + "strings" + "sync" + "sync/atomic" + "testing" + "time" +) + +func f32(v float32) *float32 { return &v } + +// The same cases as web/src/lib/player/gain.test.ts and Android's +// GainMathTest: a speaker must level a track exactly as the phone would. +func TestLeveledGainDB(t *testing.T) { + g := ReplayGain{TrackGain: f32(-6), TrackPeak: f32(1), AlbumGain: f32(-4), AlbumPeak: f32(1)} + auto := DefaultNormalizationPrefs + with := func(mode string, target int16, boost string) NormalizationPrefs { + return NormalizationPrefs{Mode: mode, TargetLUFS: target, Boost: boost} + } + cases := []struct { + name string + prefs NormalizationPrefs + g ReplayGain + asAlbum bool + want float64 + }{ + {"off", with("off", -18, "headroom"), g, false, 0}, + {"unmeasured", auto, ReplayGain{}, false, 0}, + {"track mode ignores album play", with("track", -18, "headroom"), g, true, -6}, + {"album mode", with("album", -18, "headroom"), g, false, -4}, + {"auto in album order", auto, g, true, -4}, + {"auto in a mix", auto, g, false, -6}, + {"album falls back to track", with("album", -18, "headroom"), ReplayGain{TrackGain: f32(-6), TrackPeak: f32(1)}, true, -6}, + {"louder target", with("track", -14, "headroom"), g, false, -2}, + {"headroom stops under the peak", with("track", -18, "headroom"), ReplayGain{TrackGain: f32(8), TrackPeak: f32(0.5)}, false, -1 - 20*math.Log10(0.5)}, + {"limiter lets the boost through", with("track", -18, "limiter"), ReplayGain{TrackGain: f32(8), TrackPeak: f32(0.5)}, false, 8}, + {"boost cap", with("track", -18, "limiter"), ReplayGain{TrackGain: f32(30), TrackPeak: f32(0.001)}, false, 12}, + {"cuts ignore the peak", with("track", -18, "headroom"), ReplayGain{TrackGain: f32(-9), TrackPeak: f32(1.4)}, false, -9}, + } + for _, c := range cases { + if got := LeveledGainDB(c.prefs, c.g, c.asAlbum); math.Abs(got-c.want) > 1e-4 { + t.Errorf("%s: gain = %.4f, want %.4f", c.name, got, c.want) + } + } +} + +func TestNewLeveledGain(t *testing.T) { + if g := NewLeveledGain(-6.126, "limiter"); g.CentiDB != -613 || g.Limiter { + t.Errorf("cut = %+v, want -613 without the limiter: a cut cannot raise a peak", g) + } + if g := NewLeveledGain(3.2, "limiter"); g.CentiDB != 320 || !g.Limiter { + t.Errorf("boost = %+v, want 320 with the limiter", g) + } + if g := NewLeveledGain(3.2, "headroom"); g.Limiter { + t.Errorf("headroom boost engaged the limiter: %+v", g) + } + if !NewLeveledGain(0.001, "limiter").Unity() { + t.Error("a gain that rounds to 0 should be unity") + } + for _, c := range []struct { + g LeveledGain + want bool + }{{LeveledGain{CentiDB: 1200}, true}, {LeveledGain{CentiDB: 1201}, false}, {LeveledGain{CentiDB: -6000}, true}, {LeveledGain{CentiDB: -6001}, false}} { + if c.g.Valid() != c.want { + t.Errorf("Valid(%d) = %v, want %v", c.g.CentiDB, !c.want, c.want) + } + } +} + +func TestLeveledRenderArgs(t *testing.T) { + join := func(a []string) string { return strings.Join(a, " ") } + + hires := join(leveledRenderArgs("in.flac", "out.flac", LeveledGain{CentiDB: 350, Limiter: true}, leveledSource{SampleRate: 96000, Bits: 24})) + for _, want := range []string{ + "-map_metadata -1", + "-af volume=3.50dB,alimiter=limit=0.891:level=disabled", + "-c:a flac", + "-sample_fmt s32 -bits_per_raw_sample 24", + "-ar 48000", + "-f flac -y out.flac", + } { + if !strings.Contains(hires, want) { + t.Errorf("hi-res args %q lack %q", hires, want) + } + } + + lossy := join(leveledRenderArgs("in.mp3", "out.flac", LeveledGain{CentiDB: -612}, leveledSource{SampleRate: 44100})) + if !strings.Contains(lossy, "-af volume=-6.12dB -c:a") || strings.Contains(lossy, "alimiter") { + t.Errorf("cut args %q: want a plain volume filter", lossy) + } + if !strings.Contains(lossy, "-sample_fmt s16") || strings.Contains(lossy, "-ar ") { + t.Errorf("lossy 44.1 kHz args %q: want 16-bit at the source rate", lossy) + } +} + +func TestParseLeveledProbe(t *testing.T) { + got := parseLeveledProbe("sample_rate=96000\nbits_per_sample=0\nbits_per_raw_sample=24\n") + if got != (leveledSource{SampleRate: 96000, Bits: 24}) { + t.Errorf("flac probe = %+v", got) + } + got = parseLeveledProbe("sample_rate=44100\nbits_per_sample=0\nbits_per_raw_sample=N/A\n") + if got != (leveledSource{SampleRate: 44100}) { + t.Errorf("mp3 probe = %+v", got) + } +} + +// testRenderer renders by writing size bytes, counting the renders it runs. +func testRenderer(t *testing.T, cacheMB int32, size int) (*LeveledRenderer, *atomic.Int32, chan struct{}) { + t.Helper() + settings := &LoudnessSettingsService{cur: DefaultLoudnessSettings} + settings.cur.LeveledCacheMB = cacheMB + r, err := NewLeveledRenderer(t.TempDir(), settings, slog.New(slog.NewTextHandler(io.Discard, nil))) + if err != nil { + t.Fatal(err) + } + var renders atomic.Int32 + gate := make(chan struct{}) + close(gate) // open unless a test replaces it + r.probe = func(context.Context, string) (leveledSource, error) { return leveledSource{}, nil } + r.render = func(_ context.Context, _, dst string, _ LeveledGain, _ leveledSource) error { + <-gate + renders.Add(1) + return os.WriteFile(dst, make([]byte, size), 0o600) + } + return r, &renders, gate +} + +func writeSource(t *testing.T, name string) LeveledSource { + t.Helper() + p := filepath.Join(t.TempDir(), name) + if err := os.WriteFile(p, []byte("audio"), 0o600); err != nil { + t.Fatal(err) + } + return LeveledSource{TrackID: name, Path: p, DurationMs: 1000} +} + +func TestLeveledRenderer_CoalescesAndCaches(t *testing.T) { + r, renders, _ := testRenderer(t, 2048, 10) + gate := make(chan struct{}) + r.render = func(_ context.Context, _, dst string, _ LeveledGain, _ leveledSource) error { + <-gate + renders.Add(1) + return os.WriteFile(dst, []byte("flac"), 0o600) + } + src := writeSource(t, "a") + g := LeveledGain{CentiDB: -300} + + var wg sync.WaitGroup + paths := make([]string, 5) + for i := range paths { + wg.Add(1) + go func() { + defer wg.Done() + p, err := r.Path(context.Background(), src, g) + if err != nil { + t.Errorf("Path: %v", err) + } + paths[i] = p + }() + } + time.Sleep(50 * time.Millisecond) // let every caller reach the render + close(gate) + wg.Wait() + if n := renders.Load(); n != 1 { + t.Fatalf("5 concurrent requests ran %d renders, want 1", n) + } + for _, p := range paths { + if p != paths[0] { + t.Fatalf("callers got different files: %v", paths) + } + } + + if _, err := r.Path(context.Background(), src, g); err != nil || renders.Load() != 1 { + t.Fatalf("a cached render was rendered again (renders %d, err %v)", renders.Load(), err) + } + if _, err := r.Path(context.Background(), src, LeveledGain{CentiDB: -200}); err != nil || renders.Load() != 2 { + t.Fatalf("a different gain did not render anew (renders %d, err %v)", renders.Load(), err) + } + // A replaced file is a different source: the old render must not serve it. + later := time.Now().Add(time.Hour) + if err := os.Chtimes(src.Path, later, later); err != nil { + t.Fatal(err) + } + if _, err := r.Path(context.Background(), src, g); err != nil || renders.Load() != 3 { + t.Fatalf("a changed source was served from the old render (renders %d, err %v)", renders.Load(), err) + } +} + +func TestLeveledRenderer_MissingSource(t *testing.T) { + r, _, _ := testRenderer(t, 2048, 10) + _, err := r.Path(context.Background(), LeveledSource{TrackID: "x", Path: "/nonexistent/x.flac"}, LeveledGain{CentiDB: 100}) + if !errors.Is(err, ErrLeveledSourceMissing) { + t.Fatalf("err = %v, want ErrLeveledSourceMissing", err) + } +} + +func TestLeveledRenderer_FailedRenderLeavesNothing(t *testing.T) { + r, _, _ := testRenderer(t, 2048, 10) + r.render = func(_ context.Context, _, dst string, _ LeveledGain, _ leveledSource) error { + _ = os.WriteFile(dst, []byte("half"), 0o600) + return errors.New("ffmpeg exited 1") + } + if _, err := r.Path(context.Background(), writeSource(t, "a"), LeveledGain{CentiDB: 100}); err == nil { + t.Fatal("a failed render reported success") + } + if entries, _ := os.ReadDir(r.dir); len(entries) != 0 { + t.Fatalf("a failed render left %d files behind", len(entries)) + } +} + +func TestLeveledRenderer_EvictsLeastRecentlyUsed(t *testing.T) { + // 1 MB cap, 400 KB renders: the third render pushes out the oldest. + r, _, _ := testRenderer(t, 1, 400<<10) + a, b, c := writeSource(t, "a"), writeSource(t, "b"), writeSource(t, "c") + g := LeveledGain{CentiDB: 100} + pa, _ := r.Path(context.Background(), a, g) + past := time.Now().Add(-time.Hour) + _ = os.Chtimes(pa, past, past) + pb, _ := r.Path(context.Background(), b, g) + _ = os.Chtimes(pb, past.Add(time.Minute), past.Add(time.Minute)) + // Using a again makes b the least recently used. + if _, err := r.Path(context.Background(), a, g); err != nil { + t.Fatal(err) + } + pc, _ := r.Path(context.Background(), c, g) + + var left []string + entries, _ := os.ReadDir(r.dir) + for _, e := range entries { + left = append(left, filepath.Join(r.dir, e.Name())) + } + if slices.Contains(left, pb) || !slices.Contains(left, pa) || !slices.Contains(left, pc) { + t.Fatalf("cache after eviction = %v; want a and c kept, b (least recently used) gone", left) + } +} diff --git a/internal/library/loudness_settings.go b/internal/library/loudness_settings.go index 31beca20..7f6f73d8 100644 --- a/internal/library/loudness_settings.go +++ b/internal/library/loudness_settings.go @@ -22,17 +22,27 @@ type LoudnessSettings struct { // values, so normalization keeps working for them. Enabled bool BackfillConcurrency int32 + // LeveledCacheMB caps the disk the leveled-stream renders for Sonos and + // UPnP speakers may use (#5001); the least recently played go first. + LeveledCacheMB int32 // UpdatedAt is set by the database; ignored by Set. UpdatedAt time.Time } -// DefaultLoudnessSettings mirrors migration 0065's column defaults, so a +// DefaultLoudnessSettings mirrors migrations 0065 and 0068's column defaults, so a // database that cannot be read still analyzes the way a fresh install does. var DefaultLoudnessSettings = LoudnessSettings{ Enabled: true, BackfillConcurrency: loudnessBackfillConcurrency, + LeveledCacheMB: 2048, } +// The leveled-stream cache's bounds, as migration 0068's CHECK has them. +const ( + minLeveledCacheMB = 256 + maxLeveledCacheMB = 65536 +) + // ErrLoudnessSettingOutOfRange is returned by Set for a value migration 0065's // CHECK would reject, so the API answers 400 naming the field. var ErrLoudnessSettingOutOfRange = errors.New("loudness setting out of range") @@ -78,6 +88,7 @@ func (s *LoudnessSettingsService) Set(ctx context.Context, in LoudnessSettings) row, err := dbq.New(s.pool).UpdateLoudnessSettings(ctx, dbq.UpdateLoudnessSettingsParams{ Enabled: in.Enabled, BackfillConcurrency: in.BackfillConcurrency, + LeveledCacheMb: in.LeveledCacheMB, }) if err != nil { return LoudnessSettings{}, fmt.Errorf("loudness settings: save: %w", err) @@ -94,6 +105,10 @@ func validateLoudnessSettings(in LoudnessSettings) error { return fmt.Errorf("%w: backfill_concurrency must be %d-%d", ErrLoudnessSettingOutOfRange, minBackfillConcurrency, maxBackfillConcurrency) } + if in.LeveledCacheMB < minLeveledCacheMB || in.LeveledCacheMB > maxLeveledCacheMB { + return fmt.Errorf("%w: leveled_cache_mb must be %d-%d", + ErrLoudnessSettingOutOfRange, minLeveledCacheMB, maxLeveledCacheMB) + } return nil } @@ -101,6 +116,7 @@ func loudnessSettingsFromRow(row dbq.LoudnessSetting) LoudnessSettings { return LoudnessSettings{ Enabled: row.Enabled, BackfillConcurrency: row.BackfillConcurrency, + LeveledCacheMB: row.LeveledCacheMb, UpdatedAt: row.UpdatedAt.Time, } } diff --git a/internal/library/loudness_test.go b/internal/library/loudness_test.go index d2f40704..6d4ac4da 100644 --- a/internal/library/loudness_test.go +++ b/internal/library/loudness_test.go @@ -233,14 +233,27 @@ func TestBackfillLoudnessResult_Add(t *testing.T) { } func TestValidateLoudnessSettings(t *testing.T) { - for _, n := range []int32{minBackfillConcurrency, maxBackfillConcurrency} { - if err := validateLoudnessSettings(LoudnessSettings{BackfillConcurrency: n}); err != nil { - t.Errorf("concurrency %d rejected: %v", n, err) + with := func(concurrency, cacheMB int32) LoudnessSettings { + s := DefaultLoudnessSettings + s.BackfillConcurrency, s.LeveledCacheMB = concurrency, cacheMB + return s + } + for _, s := range []LoudnessSettings{ + with(minBackfillConcurrency, minLeveledCacheMB), + with(maxBackfillConcurrency, maxLeveledCacheMB), + } { + if err := validateLoudnessSettings(s); err != nil { + t.Errorf("%+v rejected: %v", s, err) } } - for _, n := range []int32{0, maxBackfillConcurrency + 1} { - if err := validateLoudnessSettings(LoudnessSettings{BackfillConcurrency: n}); !errors.Is(err, ErrLoudnessSettingOutOfRange) { - t.Errorf("concurrency %d: err = %v, want ErrLoudnessSettingOutOfRange", n, err) + for _, s := range []LoudnessSettings{ + with(0, minLeveledCacheMB), + with(maxBackfillConcurrency+1, minLeveledCacheMB), + with(minBackfillConcurrency, minLeveledCacheMB-1), + with(minBackfillConcurrency, maxLeveledCacheMB+1), + } { + if err := validateLoudnessSettings(s); !errors.Is(err, ErrLoudnessSettingOutOfRange) { + t.Errorf("%+v: err = %v, want ErrLoudnessSettingOutOfRange", s, err) } } var nilSvc *LoudnessSettingsService diff --git a/web/src/lib/api/admin.ts b/web/src/lib/api/admin.ts index fa925e28..b257cb6d 100644 --- a/web/src/lib/api/admin.ts +++ b/web/src/lib/api/admin.ts @@ -386,6 +386,8 @@ export async function getLoudnessCoverage(): Promise { export type LoudnessSettings = { enabled: boolean; backfill_concurrency: number; + /** Disk for the leveled copies rendered for Sonos/UPnP speakers (#5001). */ + leveled_cache_mb: number; }; export async function getLoudnessSettings(): Promise { diff --git a/web/src/lib/components/LoudnessSettingsCard.svelte b/web/src/lib/components/LoudnessSettingsCard.svelte index 6003f1df..7d62066a 100644 --- a/web/src/lib/components/LoudnessSettingsCard.svelte +++ b/web/src/lib/components/LoudnessSettingsCard.svelte @@ -11,7 +11,8 @@ import { pushToast } from '$lib/stores/toast.svelte'; // Loudness analysis (#4995): the background measurement that loudness - // normalization levels playback from, its progress, and its two knobs. + // normalization levels playback from, its progress, and its knobs. The + // cache size bounds the leveled copies rendered for speakers (#5001). let saved = $state(null); let form = $state(null); @@ -26,6 +27,13 @@ form.backfill_concurrency >= 1 && form.backfill_concurrency <= 8 ); + const cacheOk = $derived( + !!form && + Number.isInteger(form.leveled_cache_mb) && + form.leveled_cache_mb >= 256 && + form.leveled_cache_mb <= 65536 + ); + const valid = $derived(concurrencyOk && cacheOk); async function load() { try { @@ -55,7 +63,7 @@ }); async function save() { - if (!form || !concurrencyOk) return; + if (!form || !valid) return; saving = true; try { saved = await updateLoudnessSettings(form); @@ -151,10 +159,26 @@ /> - {#if !concurrencyOk} -

- Files analyzed at once must be a whole number from 1 to 8. -

+ + + {#if !valid} +
+ {#if !concurrencyOk}

Files analyzed at once must be a whole number from 1 to 8.

{/if} + {#if !cacheOk}

Speaker cache must be a whole number of MB from 256 to 65536.

{/if} +
{/if}
@@ -163,7 +187,7 @@ class="rounded-md bg-action-secondary px-4 py-2 text-sm text-action-fg hover:opacity-90 focus-visible:outline focus-visible:outline-2 focus-visible:outline-accent disabled:cursor-not-allowed disabled:opacity-50" - disabled={!dirty || saving || !concurrencyOk} + disabled={!dirty || saving || !valid} onclick={save} > {saving ? 'Saving…' : 'Save'} diff --git a/web/src/lib/components/LoudnessSettingsCard.test.ts b/web/src/lib/components/LoudnessSettingsCard.test.ts index 670676ad..f1bbca8e 100644 --- a/web/src/lib/components/LoudnessSettingsCard.test.ts +++ b/web/src/lib/components/LoudnessSettingsCard.test.ts @@ -14,7 +14,7 @@ import LoudnessSettingsCard from './LoudnessSettingsCard.svelte'; import { getLoudnessCoverage, getLoudnessSettings, updateLoudnessSettings } from '$lib/api/admin'; import { pushToast } from '$lib/stores/toast.svelte'; -const base: LoudnessSettings = { enabled: true, backfill_concurrency: 2 }; +const base: LoudnessSettings = { enabled: true, backfill_concurrency: 2, leveled_cache_mb: 2048 }; const coverage: LoudnessCoverage = { total: 1200, measured: 900, @@ -68,7 +68,11 @@ describe('LoudnessSettingsCard', () => { await fireEvent.click(saveButton()); await waitFor(() => - expect(updateLoudnessSettings).toHaveBeenCalledWith({ enabled: true, backfill_concurrency: 4 }) + expect(updateLoudnessSettings).toHaveBeenCalledWith({ + enabled: true, + backfill_concurrency: 4, + leveled_cache_mb: 2048 + }) ); await waitFor(() => expect(pushToast).toHaveBeenCalledWith('Loudness analysis settings saved.')); // The gauge is read again after a save: switching analysis on or off changes it. @@ -84,6 +88,16 @@ describe('LoudnessSettingsCard', () => { expect(saveButton()).toHaveProperty('disabled', true); }); + test('the speaker cache size is bounded too', async () => { + await renderCard(); + const cache = screen.getByRole('spinbutton', { name: /speaker cache/i }); + await fireEvent.input(cache, { target: { value: '100' } }); + await waitFor(() => + expect(screen.getByTestId('settings-problems').textContent).toMatch(/256 to 65536/) + ); + expect(saveButton()).toHaveProperty('disabled', true); + }); + test('a failed load offers a retry', async () => { vi.mocked(getLoudnessSettings).mockRejectedValue(new Error('boom')); vi.mocked(getLoudnessCoverage).mockResolvedValue(coverage); diff --git a/web/src/routes/admin/admin.test.ts b/web/src/routes/admin/admin.test.ts index 04f80b6f..408a600e 100644 --- a/web/src/routes/admin/admin.test.ts +++ b/web/src/routes/admin/admin.test.ts @@ -50,7 +50,9 @@ vi.mock('$lib/api/admin', async () => { refetchMissingCovers: vi.fn().mockResolvedValue({ started: true }), researchMissingArt: vi.fn().mockResolvedValue({ version: 1 }), // LoudnessSettingsCard, rendered by the page and tested on its own. - getLoudnessSettings: vi.fn().mockResolvedValue({ enabled: true, backfill_concurrency: 2 }), + getLoudnessSettings: vi + .fn() + .mockResolvedValue({ enabled: true, backfill_concurrency: 2, leveled_cache_mb: 2048 }), updateLoudnessSettings: vi.fn(), getLoudnessCoverage: vi.fn().mockResolvedValue({ total: 0, measured: 0, silent: 0, unreadable: 0, pending: 0, enabled: true