mirror of
https://github.com/navidrome/navidrome.git
synced 2026-08-31 07:30:32 +00:00
Merge branch 'master' into feat/support-playlist-paths
This commit is contained in:
commit
b5c149285b
@ -77,6 +77,13 @@ func (s *Router) getLinkStatus(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
resp["status"] = key != ""
|
||||
linkToken, err := createLinkToken(u.ID)
|
||||
if err != nil {
|
||||
log.Error(r.Context(), "Could not create LastFM link token", "userId", u.ID, err)
|
||||
_ = rest.RespondWithError(w, http.StatusInternalServerError, err.Error())
|
||||
return
|
||||
}
|
||||
resp["linkToken"] = linkToken
|
||||
_ = rest.RespondWithJSON(w, http.StatusOK, resp)
|
||||
}
|
||||
|
||||
@ -97,11 +104,17 @@ func (s *Router) callback(w http.ResponseWriter, r *http.Request) {
|
||||
_ = rest.RespondWithError(w, http.StatusBadRequest, "token not received")
|
||||
return
|
||||
}
|
||||
uid, err := p.String("uid")
|
||||
linkToken, err := p.String("uid")
|
||||
if err != nil {
|
||||
_ = rest.RespondWithError(w, http.StatusBadRequest, "uid not received")
|
||||
return
|
||||
}
|
||||
uid, err := verifyLinkToken(linkToken)
|
||||
if err != nil {
|
||||
log.Warn(r.Context(), "Rejected LastFM callback with invalid link token", "requestId", middleware.GetReqID(r.Context()), err)
|
||||
_ = rest.RespondWithError(w, http.StatusBadRequest, "invalid link token")
|
||||
return
|
||||
}
|
||||
|
||||
// Need to add user to context, as this is a non-authenticated endpoint, so it does not
|
||||
// automatically contain any user info
|
||||
|
||||
218
adapters/lastfm/auth_router_test.go
Normal file
218
adapters/lastfm/auth_router_test.go
Normal file
@ -0,0 +1,218 @@
|
||||
package lastfm
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"time"
|
||||
|
||||
"github.com/navidrome/navidrome/core/agents"
|
||||
"github.com/navidrome/navidrome/core/auth"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/model/request"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
)
|
||||
|
||||
var _ = Describe("auth_router", func() {
|
||||
var (
|
||||
ds *tests.MockDataStore
|
||||
userProps *tests.MockedUserPropsRepo
|
||||
httpClient *tests.FakeHttpClient
|
||||
router *Router
|
||||
)
|
||||
|
||||
const (
|
||||
victimID = "victim-user-id"
|
||||
attackerID = "attacker-user-id"
|
||||
)
|
||||
|
||||
BeforeEach(func() {
|
||||
userProps = &tests.MockedUserPropsRepo{}
|
||||
ds = &tests.MockDataStore{
|
||||
MockedProperty: &tests.MockedPropertyRepo{},
|
||||
MockedUserProps: userProps,
|
||||
}
|
||||
auth.Init(ds)
|
||||
|
||||
httpClient = &tests.FakeHttpClient{}
|
||||
router = &Router{
|
||||
ds: ds,
|
||||
apiKey: "API_KEY",
|
||||
secret: "SECRET",
|
||||
sessionKeys: &agents.SessionKeys{DataStore: ds, KeyName: sessionKeyProperty},
|
||||
}
|
||||
router.client = newClient(router.apiKey, router.secret, httpClient)
|
||||
router.Handler = router.routes()
|
||||
})
|
||||
|
||||
storedSessionKey := func(userID string) string {
|
||||
key, _ := userProps.Get(userID, sessionKeyProperty)
|
||||
return key
|
||||
}
|
||||
|
||||
stubGetSessionOK := func(sessionKey string) {
|
||||
httpClient.Res = http.Response{
|
||||
Body: io.NopCloser(bytes.NewBufferString(`{"session":{"name":"Navidrome","key":"` + sessionKey + `","subscriber":0}}`)),
|
||||
StatusCode: 200,
|
||||
}
|
||||
}
|
||||
|
||||
Describe("getLinkStatus", func() {
|
||||
It("includes a signed linkToken for the authenticated user", func() {
|
||||
req := httptest.NewRequest(http.MethodGet, "/link", nil)
|
||||
ctx := request.WithUser(req.Context(), model.User{ID: victimID})
|
||||
req = req.WithContext(ctx)
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
router.getLinkStatus(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusOK))
|
||||
var body map[string]any
|
||||
Expect(json.Unmarshal(rec.Body.Bytes(), &body)).To(Succeed())
|
||||
Expect(body["apiKey"]).To(Equal("API_KEY"))
|
||||
Expect(body["status"]).To(Equal(false))
|
||||
token, ok := body["linkToken"].(string)
|
||||
Expect(ok).To(BeTrue())
|
||||
Expect(token).ToNot(BeEmpty())
|
||||
|
||||
verified, err := verifyLinkToken(token)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(verified).To(Equal(victimID))
|
||||
})
|
||||
})
|
||||
|
||||
Describe("callback", func() {
|
||||
It("stores the session key under the user encoded in the signed token", func() {
|
||||
stubGetSessionOK("LEGIT_SESSION")
|
||||
linkToken, err := createLinkToken(victimID)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/link/callback?uid="+linkToken+"&token=LASTFM_TOKEN", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
router.callback(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusOK))
|
||||
Expect(storedSessionKey(victimID)).To(Equal("LEGIT_SESSION"))
|
||||
})
|
||||
|
||||
It("rejects a raw (unsigned) uid value", func() {
|
||||
req := httptest.NewRequest(http.MethodGet, "/link/callback?uid="+victimID+"&token=LASTFM_TOKEN", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
router.callback(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusBadRequest))
|
||||
Expect(storedSessionKey(victimID)).To(BeEmpty())
|
||||
Expect(httpClient.SavedRequest).To(BeNil())
|
||||
})
|
||||
|
||||
It("rejects an expired link token", func() {
|
||||
expiredToken, err := auth.EncodeToken(map[string]any{
|
||||
"uid": victimID,
|
||||
"scope": linkTokenScope,
|
||||
"exp": time.Now().Add(-1 * time.Minute).UTC().Unix(),
|
||||
})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/link/callback?uid="+expiredToken+"&token=LASTFM_TOKEN", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
router.callback(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusBadRequest))
|
||||
Expect(storedSessionKey(victimID)).To(BeEmpty())
|
||||
Expect(httpClient.SavedRequest).To(BeNil())
|
||||
})
|
||||
|
||||
It("rejects a token with the wrong scope (e.g. a regular session JWT)", func() {
|
||||
sessionJWT, err := auth.CreateToken(&model.User{ID: attackerID, UserName: "attacker"})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/link/callback?uid="+sessionJWT+"&token=LASTFM_TOKEN", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
router.callback(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusBadRequest))
|
||||
Expect(storedSessionKey(attackerID)).To(BeEmpty())
|
||||
Expect(httpClient.SavedRequest).To(BeNil())
|
||||
})
|
||||
|
||||
It("writes only under the user encoded in the token, regardless of query manipulation", func() {
|
||||
// An attacker holds a legitimate link token for their own account.
|
||||
// They attempt to call the callback hoping to overwrite the victim's
|
||||
// session key — but the handler must derive the user ID from the
|
||||
// signed token, not from any other input.
|
||||
stubGetSessionOK("ATTACKER_SESSION")
|
||||
attackerToken, err := createLinkToken(attackerID)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/link/callback?uid="+attackerToken+"&token=LASTFM_TOKEN&user="+victimID, nil)
|
||||
rec := httptest.NewRecorder()
|
||||
router.callback(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusOK))
|
||||
Expect(storedSessionKey(attackerID)).To(Equal("ATTACKER_SESSION"))
|
||||
Expect(storedSessionKey(victimID)).To(BeEmpty())
|
||||
})
|
||||
|
||||
It("returns 400 when uid is missing", func() {
|
||||
req := httptest.NewRequest(http.MethodGet, "/link/callback?token=LASTFM_TOKEN", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
router.callback(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusBadRequest))
|
||||
})
|
||||
|
||||
It("returns 400 when token is missing", func() {
|
||||
linkToken, err := createLinkToken(victimID)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/link/callback?uid="+linkToken, nil)
|
||||
rec := httptest.NewRecorder()
|
||||
router.callback(rec, req)
|
||||
|
||||
Expect(rec.Code).To(Equal(http.StatusBadRequest))
|
||||
})
|
||||
})
|
||||
|
||||
Describe("link token helpers", func() {
|
||||
It("round-trips a freshly issued token", func() {
|
||||
token, err := createLinkToken(victimID)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
uid, err := verifyLinkToken(token)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(uid).To(Equal(victimID))
|
||||
})
|
||||
|
||||
It("rejects garbage", func() {
|
||||
_, err := verifyLinkToken("not-a-jwt")
|
||||
Expect(err).To(HaveOccurred())
|
||||
})
|
||||
|
||||
It("rejects a token whose scope claim is wrong", func() {
|
||||
wrongScopeToken, err := auth.EncodeToken(map[string]any{
|
||||
"uid": victimID,
|
||||
"scope": "some-other-scope",
|
||||
"exp": time.Now().Add(linkTokenTTL).UTC().Unix(),
|
||||
})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
_, err = verifyLinkToken(wrongScopeToken)
|
||||
Expect(err).To(MatchError("invalid link token scope"))
|
||||
})
|
||||
|
||||
It("rejects a scoped token that has no expiration", func() {
|
||||
nonExpiringToken, err := auth.EncodeToken(map[string]any{
|
||||
"uid": victimID,
|
||||
"scope": linkTokenScope,
|
||||
})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
_, err = verifyLinkToken(nonExpiringToken)
|
||||
Expect(err).To(MatchError("link token missing expiration"))
|
||||
})
|
||||
})
|
||||
})
|
||||
50
adapters/lastfm/link_token.go
Normal file
50
adapters/lastfm/link_token.go
Normal file
@ -0,0 +1,50 @@
|
||||
package lastfm
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"time"
|
||||
|
||||
"github.com/navidrome/navidrome/core/auth"
|
||||
)
|
||||
|
||||
const (
|
||||
linkTokenScope = "lastfm-link"
|
||||
linkTokenTTL = 5 * time.Minute
|
||||
)
|
||||
|
||||
// createLinkToken issues a signed token binding the Last.fm callback to the
|
||||
// user who initiated the OAuth flow. It travels back through Last.fm via the
|
||||
// `cb` URL in place of the previously-trusted raw `uid` query parameter.
|
||||
func createLinkToken(userID string) (string, error) {
|
||||
claims := map[string]any{
|
||||
"uid": userID,
|
||||
"scope": linkTokenScope,
|
||||
"exp": time.Now().Add(linkTokenTTL).UTC().Unix(),
|
||||
}
|
||||
return auth.EncodeToken(claims)
|
||||
}
|
||||
|
||||
// verifyLinkToken validates a signed link token and returns the encoded user ID.
|
||||
// It enforces both the signature/expiry (via the underlying JWT verifier) and a
|
||||
// dedicated scope claim, preventing tokens minted for other purposes (e.g. a
|
||||
// regular session JWT) from being accepted here.
|
||||
func verifyLinkToken(tokenStr string) (string, error) {
|
||||
token, err := auth.DecodeAndVerifyToken(tokenStr)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
// jwtauth treats a token without `exp` as non-expiring; require it
|
||||
// explicitly so an accidental regression cannot mint permanent tokens.
|
||||
if exp, ok := token.Expiration(); !ok || exp.IsZero() {
|
||||
return "", errors.New("link token missing expiration")
|
||||
}
|
||||
var scope string
|
||||
if err := token.Get("scope", &scope); err != nil || scope != linkTokenScope {
|
||||
return "", errors.New("invalid link token scope")
|
||||
}
|
||||
var uid string
|
||||
if err := token.Get("uid", &uid); err != nil || uid == "" {
|
||||
return "", errors.New("invalid link token user ID")
|
||||
}
|
||||
return uid, nil
|
||||
}
|
||||
@ -32,17 +32,17 @@ var inspectCmd = &cobra.Command{
|
||||
},
|
||||
}
|
||||
|
||||
var marshalers = map[string]func(interface{}) ([]byte, error){
|
||||
var marshalers = map[string]func(any) ([]byte, error){
|
||||
"pretty": prettyMarshal,
|
||||
"toml": toml.Marshal,
|
||||
"yaml": yaml.Marshal,
|
||||
"json": json.Marshal,
|
||||
"jsonindent": func(v interface{}) ([]byte, error) {
|
||||
"jsonindent": func(v any) ([]byte, error) {
|
||||
return json.MarshalIndent(v, "", " ")
|
||||
},
|
||||
}
|
||||
|
||||
func prettyMarshal(v interface{}) ([]byte, error) {
|
||||
func prettyMarshal(v any) ([]byte, error) {
|
||||
out := v.([]core.InspectOutput)
|
||||
var res strings.Builder
|
||||
for i := range out {
|
||||
|
||||
@ -47,7 +47,6 @@ type configOptions struct {
|
||||
UIWelcomeMessage string
|
||||
MaxSidebarPlaylists int
|
||||
EnableTranscodingConfig bool
|
||||
EnableTranscodingCancellation bool
|
||||
EnableDownloads bool
|
||||
EnableExternalServices bool
|
||||
EnableM3UExternalAlbumArt bool
|
||||
@ -113,6 +112,7 @@ type configOptions struct {
|
||||
PID pidOptions `json:",omitzero"`
|
||||
Inspect inspectOptions `json:",omitzero"`
|
||||
Subsonic subsonicOptions `json:",omitzero"`
|
||||
Transcoding transcodingOptions `json:",omitzero"`
|
||||
LastFM lastfmOptions `json:",omitzero"`
|
||||
Deezer deezerOptions `json:",omitzero"`
|
||||
ListenBrainz listenBrainzOptions `json:",omitzero"`
|
||||
@ -165,6 +165,12 @@ type scannerOptions struct {
|
||||
PurgeMissing string // Values: "never", "always", "full"
|
||||
}
|
||||
|
||||
type transcodingOptions struct {
|
||||
MaxConcurrent int
|
||||
MaxConcurrentPerUser int
|
||||
EnableCancellation bool
|
||||
}
|
||||
|
||||
type subsonicOptions struct {
|
||||
AppendSubtitle bool
|
||||
AppendAlbumVersion bool
|
||||
@ -324,6 +330,7 @@ func Load(noConfigDump bool) {
|
||||
mapDeprecatedOption("HTTPSecurityHeaders.CustomFrameOptionsValue", "HTTPHeaders.FrameOptions")
|
||||
mapDeprecatedOption("CoverJpegQuality", "CoverArtQuality")
|
||||
mapDeprecatedOption("SimilarSongsMatchThreshold", "Matcher.FuzzyThreshold")
|
||||
mapDeprecatedOption("EnableTranscodingCancellation", "Transcoding.EnableCancellation")
|
||||
|
||||
err := viper.Unmarshal(&Server, viper.DecodeHook(
|
||||
mapstructure.ComposeDecodeHookFunc(
|
||||
@ -449,6 +456,7 @@ func Load(noConfigDump bool) {
|
||||
logDeprecatedOptions("HTTPSecurityHeaders.CustomFrameOptionsValue", "HTTPHeaders.FrameOptions")
|
||||
logDeprecatedOptions("CoverJpegQuality", "CoverArtQuality")
|
||||
logDeprecatedOptions("SimilarSongsMatchThreshold", "Matcher.FuzzyThreshold")
|
||||
logDeprecatedOptions("EnableTranscodingCancellation", "Transcoding.EnableCancellation")
|
||||
|
||||
// Removed options
|
||||
logRemovedOptions("Spotify.ID", "Spotify.Secret")
|
||||
@ -737,7 +745,6 @@ func setViperDefaults() {
|
||||
viper.SetDefault("uiwelcomemessage", "")
|
||||
viper.SetDefault("maxsidebarplaylists", consts.DefaultMaxSidebarPlaylists)
|
||||
viper.SetDefault("enabletranscodingconfig", false)
|
||||
viper.SetDefault("enabletranscodingcancellation", false)
|
||||
viper.SetDefault("transcodingcachesize", "100MB")
|
||||
viper.SetDefault("imagecachesize", "100MB")
|
||||
viper.SetDefault("albumplaycountmode", consts.AlbumPlayCountModeAbsolute)
|
||||
@ -822,6 +829,9 @@ func setViperDefaults() {
|
||||
viper.SetDefault("subsonic.enableaveragerating", true)
|
||||
viper.SetDefault("subsonic.legacyclients", "DSub")
|
||||
viper.SetDefault("subsonic.minimalclients", "SubMusic")
|
||||
viper.SetDefault("transcoding.maxconcurrent", 0)
|
||||
viper.SetDefault("transcoding.maxconcurrentperuser", 0)
|
||||
viper.SetDefault("transcoding.enablecancellation", false)
|
||||
viper.SetDefault("agents", "deezer,lastfm,listenbrainz")
|
||||
viper.SetDefault("lastfm.enabled", true)
|
||||
viper.SetDefault("lastfm.language", consts.DefaultInfoLanguage)
|
||||
|
||||
47
conf/dir.go
47
conf/dir.go
@ -1,20 +1,20 @@
|
||||
package conf
|
||||
|
||||
import (
|
||||
"cmp"
|
||||
"fmt"
|
||||
"os"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// Dir wraps a directory path and lazily creates the directory on first use.
|
||||
// The directory is created at most once; if creation fails, the error is
|
||||
// permanently cached (sync.Once semantics). Dir is not safe for mutation
|
||||
// after Path() has been called.
|
||||
// Dir wraps a directory path and creates the directory on demand. Dir is a
|
||||
// plain value type — safe to copy, compare, and print via reflection-based
|
||||
// formatters (pretty.Sprintf("%# v", ...)) without any concurrency hazards.
|
||||
// Directory creation is delegated to os.MkdirAll on every Path() call;
|
||||
// MkdirAll is idempotent, so repeated calls cost one stat syscall when the
|
||||
// directory already exists.
|
||||
type Dir struct {
|
||||
path string
|
||||
perm os.FileMode
|
||||
once sync.Once
|
||||
err error
|
||||
}
|
||||
|
||||
// NewDir creates a new Dir with the given path and default permissions (os.ModePerm).
|
||||
@ -23,31 +23,32 @@ func NewDir(path string) Dir {
|
||||
}
|
||||
|
||||
// NewDirWithPerm creates a new Dir with the given path and permissions.
|
||||
// A perm of 0 is treated as "default" and resolves to os.ModePerm at
|
||||
// directory-creation time; pass an explicit non-zero mode to constrain the
|
||||
// permissions.
|
||||
func NewDirWithPerm(path string, perm os.FileMode) Dir {
|
||||
return Dir{path: path, perm: perm}
|
||||
}
|
||||
|
||||
// String returns the raw path without creating the directory. Satisfies fmt.Stringer.
|
||||
func (d *Dir) String() string {
|
||||
func (d Dir) String() string {
|
||||
return d.path
|
||||
}
|
||||
|
||||
// Path creates the directory on first call (via sync.Once) and returns the path.
|
||||
func (d *Dir) Path() (string, error) {
|
||||
d.once.Do(func() {
|
||||
if d.path == "" {
|
||||
return
|
||||
}
|
||||
d.err = os.MkdirAll(d.path, d.perm)
|
||||
if d.err != nil {
|
||||
d.err = fmt.Errorf("creating directory %q: %w", d.path, d.err)
|
||||
}
|
||||
})
|
||||
return d.path, d.err
|
||||
// Path ensures the directory exists and returns its path. Safe to call
|
||||
// repeatedly; an empty path is returned as-is with no error.
|
||||
func (d Dir) Path() (string, error) {
|
||||
if d.path == "" {
|
||||
return "", nil
|
||||
}
|
||||
if err := os.MkdirAll(d.path, cmp.Or(d.perm, os.ModePerm)); err != nil {
|
||||
return d.path, fmt.Errorf("creating directory %q: %w", d.path, err)
|
||||
}
|
||||
return d.path, nil
|
||||
}
|
||||
|
||||
// MustPath calls Path() and calls logFatal on error.
|
||||
func (d *Dir) MustPath() string {
|
||||
func (d Dir) MustPath() string {
|
||||
path, err := d.Path()
|
||||
if err != nil {
|
||||
logFatal("creating directory:", err)
|
||||
@ -57,12 +58,12 @@ func (d *Dir) MustPath() string {
|
||||
|
||||
// GoString implements fmt.GoStringer so that %#v (used by pretty.Sprintf)
|
||||
// prints the path string instead of the internal struct fields.
|
||||
func (d Dir) GoString() string { //nolint:govet // uses a value receiver so Dir values satisfy GoStringer
|
||||
func (d Dir) GoString() string {
|
||||
return fmt.Sprintf("%q", d.path)
|
||||
}
|
||||
|
||||
// MarshalText returns the raw path bytes. No side effects.
|
||||
func (d *Dir) MarshalText() ([]byte, error) {
|
||||
func (d Dir) MarshalText() ([]byte, error) {
|
||||
return []byte(d.path), nil
|
||||
}
|
||||
|
||||
|
||||
@ -2,7 +2,9 @@ package conf_test
|
||||
|
||||
import (
|
||||
"os"
|
||||
"sync"
|
||||
|
||||
"github.com/kr/pretty"
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
@ -35,9 +37,9 @@ var _ = Describe("Dir", func() {
|
||||
Expect(target).To(BeADirectory())
|
||||
})
|
||||
|
||||
It("returns the same result on subsequent calls (sync.Once)", func() {
|
||||
It("is idempotent on subsequent calls", func() {
|
||||
dir := GinkgoT().TempDir()
|
||||
target := dir + "/once"
|
||||
target := dir + "/idempotent"
|
||||
d := conf.NewDir(target)
|
||||
|
||||
path1, err1 := d.Path()
|
||||
@ -45,6 +47,7 @@ var _ = Describe("Dir", func() {
|
||||
Expect(err1).ToNot(HaveOccurred())
|
||||
Expect(err2).ToNot(HaveOccurred())
|
||||
Expect(path1).To(Equal(path2))
|
||||
Expect(target).To(BeADirectory())
|
||||
})
|
||||
|
||||
It("returns an error when directory cannot be created", func() {
|
||||
@ -124,4 +127,38 @@ var _ = Describe("Dir", func() {
|
||||
Expect(d2.String()).To(Equal(d1.String()))
|
||||
})
|
||||
})
|
||||
|
||||
Describe("GoString", func() {
|
||||
// Regression: pretty.Sprintf("%# v", ...) is used by the
|
||||
// configuration dump. It must render Dir as a quoted path via
|
||||
// GoString, not dump the internal struct fields.
|
||||
It("renders Dir as a quoted path under pretty.Sprintf", func() {
|
||||
type host struct {
|
||||
DataFolder conf.Dir
|
||||
}
|
||||
h := host{DataFolder: conf.NewDir("./data")}
|
||||
out := pretty.Sprintf("%# v", h)
|
||||
Expect(out).To(ContainSubstring(`DataFolder: "./data"`))
|
||||
Expect(out).ToNot(ContainSubstring("perm:"))
|
||||
Expect(out).ToNot(ContainSubstring("path:"))
|
||||
})
|
||||
|
||||
It("is safe to copy and use concurrently", func() {
|
||||
// Regression for the Windows "sync: unlock of unlocked mutex"
|
||||
// crash that was caused by copying a Dir embedding sync.Once.
|
||||
// Dir is a plain value type now, but keep the concurrent stress
|
||||
// test to lock in the property.
|
||||
dir := GinkgoT().TempDir()
|
||||
d := conf.NewDir(dir + "/race")
|
||||
var wg sync.WaitGroup
|
||||
for range 10 {
|
||||
wg.Go(func() {
|
||||
copy1 := d
|
||||
_ = pretty.Sprintf("%# v", copy1)
|
||||
_, _ = copy1.Path()
|
||||
})
|
||||
}
|
||||
wg.Wait()
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
@ -3,6 +3,7 @@ package core
|
||||
import (
|
||||
"archive/zip"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
@ -60,7 +61,15 @@ func (a *archiver) zipAlbums(ctx context.Context, id string, format string, bitr
|
||||
"format", format, "bitrate", bitrate, "isMultiDisc", isMultiDisc, "numTracks", len(album))
|
||||
for _, mf := range album {
|
||||
file := a.albumFilename(mf, format, isMultiDisc)
|
||||
_ = a.addFileToZip(ctx, z, mf, format, bitrate, file)
|
||||
if addErr := a.addFileToZip(ctx, z, mf, format, bitrate, file); errors.Is(addErr, stream.ErrTooManyTranscodes) {
|
||||
// Stop iterating: continuing would just rack up more
|
||||
// rejections from the limiter. Close finalises whatever
|
||||
// tracks were already written; the rejected one is not
|
||||
// present in the archive (addFileToZip aborts before
|
||||
// writing its entry header).
|
||||
_ = z.Close()
|
||||
return addErr
|
||||
}
|
||||
}
|
||||
}
|
||||
err = z.Close()
|
||||
@ -120,7 +129,12 @@ func (a *archiver) zipMediaFiles(ctx context.Context, id, name string, format st
|
||||
zippedMfs := make(model.MediaFiles, len(mfs))
|
||||
for idx, mf := range mfs {
|
||||
file := a.playlistFilename(mf, format, idx)
|
||||
_ = a.addFileToZip(ctx, z, mf, format, bitrate, file)
|
||||
if addErr := a.addFileToZip(ctx, z, mf, format, bitrate, file); errors.Is(addErr, stream.ErrTooManyTranscodes) {
|
||||
// Abort the whole archive: continuing would silently emit
|
||||
// empty zip entries since the headers are already written.
|
||||
_ = z.Close()
|
||||
return addErr
|
||||
}
|
||||
mf.Path = file
|
||||
zippedMfs[idx] = mf
|
||||
}
|
||||
@ -162,6 +176,27 @@ func (a *archiver) playlistFilename(mf model.MediaFile, format string, idx int)
|
||||
|
||||
func (a *archiver) addFileToZip(ctx context.Context, z *zip.Writer, mf model.MediaFile, format string, bitrate int, filename string) error {
|
||||
path := mf.AbsolutePath()
|
||||
|
||||
// Open the source before writing the zip entry header so a rejection
|
||||
// (limiter, missing file, etc.) does not leave an empty entry in the
|
||||
// archive.
|
||||
var r io.ReadCloser
|
||||
var err error
|
||||
if format != "raw" && format != "" {
|
||||
r, err = a.ms.NewStream(ctx, &mf, stream.Request{Format: format, BitRate: bitrate})
|
||||
} else {
|
||||
r, err = os.Open(path)
|
||||
}
|
||||
if err != nil {
|
||||
log.Error(ctx, "Error opening file for zipping", "file", path, "format", format, err)
|
||||
return err
|
||||
}
|
||||
defer func() {
|
||||
if err := r.Close(); err != nil && log.IsGreaterOrEqualTo(log.LevelDebug) {
|
||||
log.Error(ctx, "Error closing stream", "id", mf.ID, "file", path, err)
|
||||
}
|
||||
}()
|
||||
|
||||
w, err := z.CreateHeader(&zip.FileHeader{
|
||||
Name: filename,
|
||||
Modified: mf.UpdatedAt,
|
||||
@ -172,23 +207,6 @@ func (a *archiver) addFileToZip(ctx context.Context, z *zip.Writer, mf model.Med
|
||||
return err
|
||||
}
|
||||
|
||||
var r io.ReadCloser
|
||||
if format != "raw" && format != "" {
|
||||
r, err = a.ms.NewStream(ctx, &mf, stream.Request{Format: format, BitRate: bitrate})
|
||||
} else {
|
||||
r, err = os.Open(path)
|
||||
}
|
||||
if err != nil {
|
||||
log.Error(ctx, "Error opening file for zipping", "file", path, "format", format, err)
|
||||
return err
|
||||
}
|
||||
|
||||
defer func() {
|
||||
if err := r.Close(); err != nil && log.IsGreaterOrEqualTo(log.LevelDebug) {
|
||||
log.Error(ctx, "Error closing stream", "id", mf.ID, "file", path, err)
|
||||
}
|
||||
}()
|
||||
|
||||
_, err = io.Copy(w, r)
|
||||
if err != nil {
|
||||
log.Error(ctx, "Error zipping file", "file", path, err)
|
||||
|
||||
@ -89,6 +89,32 @@ var _ = Describe("Archiver", func() {
|
||||
})
|
||||
})
|
||||
|
||||
Context("when the transcode limiter rejects a file", func() {
|
||||
It("aborts the archive instead of continuing with empty entries", func() {
|
||||
mfs := model.MediaFiles{
|
||||
{Path: "test_data/01 - track1.mp3", Suffix: "mp3", AlbumID: "1", Album: "Album", DiscNumber: 1},
|
||||
{Path: "test_data/02 - track2.mp3", Suffix: "mp3", AlbumID: "1", Album: "Album", DiscNumber: 1},
|
||||
}
|
||||
|
||||
mfRepo := &mockMediaFileRepository{}
|
||||
mfRepo.On("GetAll", []model.QueryOptions{{
|
||||
Filters: squirrel.Eq{"album_id": "1"},
|
||||
Sort: "album",
|
||||
}}).Return(mfs, nil)
|
||||
ds.On("MediaFile", mock.Anything).Return(mfRepo)
|
||||
|
||||
ms.On("NewStream", mock.Anything, mock.Anything, stream.Request{Format: "mp3", BitRate: 128}).
|
||||
Return(nil, stream.ErrTooManyTranscodes).Once()
|
||||
|
||||
out := new(bytes.Buffer)
|
||||
err := arch.ZipAlbum(context.Background(), "1", "mp3", 128, out)
|
||||
Expect(err).To(MatchError(stream.ErrTooManyTranscodes))
|
||||
// NewStream should only have been called once: the loop must bail
|
||||
// out on the rejection instead of trying every remaining track.
|
||||
ms.AssertNumberOfCalls(GinkgoT(), "NewStream", 1)
|
||||
})
|
||||
})
|
||||
|
||||
Context("ZipShare", func() {
|
||||
It("zips a share correctly", func() {
|
||||
mfs := model.MediaFiles{
|
||||
|
||||
@ -169,7 +169,7 @@ func BenchmarkArtworkGetE2EConcurrent(b *testing.B) {
|
||||
for i := 0; i < b.N; i++ {
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(n)
|
||||
for g := 0; g < n; g++ {
|
||||
for range n {
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
r, _, err := aw.Get(context.Background(), artID, 300, true)
|
||||
|
||||
@ -35,8 +35,8 @@ func generatePNG(t testing.TB, width, height int) []byte {
|
||||
// generateGradientImage creates an RGBA image with a diagonal gradient pattern.
|
||||
func generateGradientImage(width, height int) *image.RGBA {
|
||||
img := image.NewRGBA(image.Rect(0, 0, width, height))
|
||||
for y := 0; y < height; y++ {
|
||||
for x := 0; x < width; x++ {
|
||||
for y := range height {
|
||||
for x := range width {
|
||||
r := uint8((x * 255) / width)
|
||||
g := uint8((y * 255) / height)
|
||||
b := uint8(((x + y) * 255) / (width + height))
|
||||
|
||||
@ -496,8 +496,8 @@ func createFFmpegCommand(cmd, path string, maxBitRate, offset int) []string {
|
||||
// Pre-input seeking: ffmpeg seeks at the demuxer level (fast)
|
||||
// instead of decoding all frames up to the offset (slow).
|
||||
insertAt := len(args)
|
||||
for i := len(args) - 1; i >= 0; i-- {
|
||||
if args[i] == "-i" {
|
||||
for i, arg := range slices.Backward(args) {
|
||||
if arg == "-i" {
|
||||
insertAt = i
|
||||
break
|
||||
}
|
||||
|
||||
@ -4,11 +4,13 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
"reflect"
|
||||
"strings"
|
||||
|
||||
"github.com/deluan/rest"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/model/criteria"
|
||||
"github.com/navidrome/navidrome/model/request"
|
||||
"github.com/navidrome/navidrome/utils/slice"
|
||||
)
|
||||
|
||||
// --- REST adapter (follows Share/Library pattern) ---
|
||||
@ -34,8 +36,8 @@ func (r *playlistRepositoryWrapper) Save(entity any) (string, error) {
|
||||
return r.service.savePlaylist(r.ctx, entity.(*model.Playlist))
|
||||
}
|
||||
|
||||
func (r *playlistRepositoryWrapper) Update(id string, entity any, _ ...string) error {
|
||||
return r.service.updatePlaylistEntity(r.ctx, id, entity.(*model.Playlist))
|
||||
func (r *playlistRepositoryWrapper) Update(id string, entity any, cols ...string) error {
|
||||
return r.service.updatePlaylistEntity(r.ctx, id, entity.(*model.Playlist), cols...)
|
||||
}
|
||||
|
||||
func (r *playlistRepositoryWrapper) Delete(id string) error {
|
||||
@ -79,7 +81,15 @@ func (s *playlists) savePlaylist(ctx context.Context, pls *model.Playlist) (stri
|
||||
|
||||
// updatePlaylistEntity updates playlist metadata with permission checks.
|
||||
// Used by the REST API wrapper.
|
||||
func (s *playlists) updatePlaylistEntity(ctx context.Context, id string, entity *model.Playlist) error {
|
||||
//
|
||||
// cols names the fields the client actually sent in the JSON body (extracted by
|
||||
// rest.Put). When non-empty, fields outside cols are not considered changed and
|
||||
// are left untouched — this prevents partial requests like bulk "Make Public"
|
||||
// (body: {"public": true}) from wiping fields that just happen to be zero in
|
||||
// the deserialized entity (see issue #5541). An empty cols means "treat the
|
||||
// entity as a complete record" — preserved for callers that don't use the REST
|
||||
// wrapper.
|
||||
func (s *playlists) updatePlaylistEntity(ctx context.Context, id string, entity *model.Playlist, cols ...string) error {
|
||||
current, err := s.checkWritable(ctx, id)
|
||||
if err != nil {
|
||||
switch {
|
||||
@ -91,41 +101,92 @@ func (s *playlists) updatePlaylistEntity(ctx context.Context, id string, entity
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
sent := sentFields(cols)
|
||||
|
||||
usr, _ := request.UserFrom(ctx)
|
||||
if !usr.IsAdmin && entity.OwnerID != "" && entity.OwnerID != current.OwnerID {
|
||||
ownerChanged := sent("ownerId") && entity.OwnerID != "" && entity.OwnerID != current.OwnerID
|
||||
if !usr.IsAdmin && ownerChanged {
|
||||
return rest.ErrPermissionDenied
|
||||
}
|
||||
|
||||
contentChanged := entity.Name != current.Name ||
|
||||
entity.Comment != current.Comment ||
|
||||
(entity.OwnerID != "" && entity.OwnerID != current.OwnerID) ||
|
||||
!rulesEqual(current.Rules, entity.Rules)
|
||||
nameChanged := sent("name") && entity.Name != current.Name
|
||||
commentChanged := sent("comment") && entity.Comment != current.Comment
|
||||
rulesChanged := sent("rules") && !rulesEqual(current.Rules, entity.Rules)
|
||||
|
||||
if contentChanged {
|
||||
if entity.OwnerID != "" {
|
||||
current.OwnerID = entity.OwnerID
|
||||
}
|
||||
if nameChanged || commentChanged || ownerChanged || rulesChanged {
|
||||
return s.applyContentUpdate(ctx, current, entity, sent,
|
||||
nameChanged, commentChanged, ownerChanged, rulesChanged)
|
||||
}
|
||||
return s.applyFlagsOnly(ctx, current, entity, sent)
|
||||
}
|
||||
|
||||
// applyContentUpdate handles updates that change at least one of name/comment/
|
||||
// owner/rules. It goes through updateMetadata, which always bumps updatedAt
|
||||
// (invalidating cached cover-art URLs). namePtr/commentPtr are nil when the
|
||||
// field is absent from the request OR present-but-unchanged (so updateMetadata
|
||||
// skips them); publicPtr is nil only when public is absent from the request
|
||||
// (an idempotent public value is still forwarded).
|
||||
func (s *playlists) applyContentUpdate(ctx context.Context, current, entity *model.Playlist,
|
||||
sent func(string) bool, nameChanged, commentChanged, ownerChanged, rulesChanged bool,
|
||||
) error {
|
||||
if ownerChanged {
|
||||
current.OwnerID = entity.OwnerID
|
||||
}
|
||||
if rulesChanged {
|
||||
current.Rules = entity.Rules
|
||||
if current.Path != "" && current.Sync != entity.Sync {
|
||||
current.Sync = entity.Sync
|
||||
}
|
||||
return s.updateMetadata(ctx, s.ds, current, &entity.Name, &entity.Comment, &entity.Public)
|
||||
}
|
||||
|
||||
// Only sync/public changed — skip updatedAt so cover art URLs stay stable
|
||||
var cols []string
|
||||
if current.Path != "" && current.Sync != entity.Sync {
|
||||
if sent("sync") && current.Path != "" && current.Sync != entity.Sync {
|
||||
current.Sync = entity.Sync
|
||||
cols = append(cols, "sync")
|
||||
}
|
||||
if current.Public != entity.Public {
|
||||
var namePtr, commentPtr *string
|
||||
var publicPtr *bool
|
||||
if nameChanged {
|
||||
namePtr = &entity.Name
|
||||
}
|
||||
if commentChanged {
|
||||
commentPtr = &entity.Comment
|
||||
}
|
||||
if sent("public") {
|
||||
publicPtr = &entity.Public
|
||||
}
|
||||
return s.updateMetadata(ctx, s.ds, current, namePtr, commentPtr, publicPtr)
|
||||
}
|
||||
|
||||
// applyFlagsOnly handles updates that only toggle sync/public — skips
|
||||
// updatedAt so cover art URLs stay stable.
|
||||
func (s *playlists) applyFlagsOnly(ctx context.Context, current, entity *model.Playlist,
|
||||
sent func(string) bool,
|
||||
) error {
|
||||
var updateCols []string
|
||||
if sent("sync") && current.Path != "" && current.Sync != entity.Sync {
|
||||
current.Sync = entity.Sync
|
||||
updateCols = append(updateCols, "sync")
|
||||
}
|
||||
if sent("public") && current.Public != entity.Public {
|
||||
current.Public = entity.Public
|
||||
cols = append(cols, "public")
|
||||
updateCols = append(updateCols, "public")
|
||||
}
|
||||
if len(cols) == 0 {
|
||||
if len(updateCols) == 0 {
|
||||
return nil
|
||||
}
|
||||
return s.ds.Playlist(ctx).Put(current, cols...)
|
||||
return s.ds.Playlist(ctx).Put(current, updateCols...)
|
||||
}
|
||||
|
||||
// sentFields returns a predicate that reports whether a JSON field was present
|
||||
// in the request body. Matching is case-insensitive to mirror Go's json
|
||||
// decoder, which populates struct fields from case-variant keys like
|
||||
// {"Name":"x"} or {"OWNERID":"y"}. An empty cols list means "treat the entity
|
||||
// as a full record" — every field is considered sent.
|
||||
func sentFields(cols []string) func(string) bool {
|
||||
if len(cols) == 0 {
|
||||
return func(string) bool { return true }
|
||||
}
|
||||
set := slice.ToMap(cols, func(c string) (string, struct{}) { return strings.ToLower(c), struct{}{} })
|
||||
return func(field string) bool {
|
||||
_, ok := set[strings.ToLower(field)]
|
||||
return ok
|
||||
}
|
||||
}
|
||||
|
||||
func rulesEqual(a, b *criteria.Criteria) bool {
|
||||
|
||||
@ -125,6 +125,25 @@ var _ = Describe("REST Adapter", func() {
|
||||
Expect(err).To(Equal(rest.ErrPermissionDenied))
|
||||
})
|
||||
|
||||
DescribeTable("denies regular user from changing ownership under any case-variant JSON key",
|
||||
func(colName string) {
|
||||
// rest.Put's field-name extraction is case-sensitive, but Go's
|
||||
// json decoder is case-insensitive on struct fields, so any
|
||||
// {"OwnerId":"x"} / {"OWNERID":"x"} / {"ownerid":"x"} populates
|
||||
// entity.OwnerID. sentFields normalizes both sides so the
|
||||
// permission gate fires regardless of casing.
|
||||
ctx = request.WithUser(ctx, model.User{ID: "user-1", IsAdmin: false})
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
pls := &model.Playlist{OwnerID: "other-user"}
|
||||
err := repo.Update("pls-1", pls, colName)
|
||||
Expect(err).To(Equal(rest.ErrPermissionDenied))
|
||||
},
|
||||
Entry("canonical camelCase", "ownerId"),
|
||||
Entry("PascalCase", "OwnerId"),
|
||||
Entry("all upper", "OWNERID"),
|
||||
Entry("all lower", "ownerid"),
|
||||
)
|
||||
|
||||
It("updates smart playlist rules", func() {
|
||||
mockPlsRepo.Data["smart-1"] = &model.Playlist{
|
||||
ID: "smart-1",
|
||||
@ -218,6 +237,156 @@ var _ = Describe("REST Adapter", func() {
|
||||
err := repo.Update("nonexistent", pls)
|
||||
Expect(err).To(Equal(rest.ErrNotFound))
|
||||
})
|
||||
|
||||
// Regression tests for #5541: partial REST updates (e.g. bulk "Make Public")
|
||||
// must only touch the fields the client actually sent. The cols list from
|
||||
// rest.Put names those fields; fields outside it must be left alone, even
|
||||
// when the deserialized entity has zero values for them.
|
||||
Context("with partial updates (cols)", func() {
|
||||
BeforeEach(func() {
|
||||
ctx = request.WithUser(ctx, model.User{ID: "user-1", IsAdmin: false})
|
||||
mockPlsRepo.Data["partial"] = &model.Playlist{
|
||||
ID: "partial",
|
||||
Name: "Original Name",
|
||||
Comment: "Original comment",
|
||||
OwnerID: "user-1",
|
||||
Public: false,
|
||||
}
|
||||
})
|
||||
|
||||
It("preserves name and comment when only public is sent (bulk Make Public)", func() {
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
err := repo.Update("partial", &model.Playlist{Public: true}, "public")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Original Name"))
|
||||
Expect(mockPlsRepo.Last.Comment).To(Equal("Original comment"))
|
||||
Expect(mockPlsRepo.Last.Public).To(BeTrue())
|
||||
})
|
||||
|
||||
It("preserves name when only sync is sent for a file-backed playlist", func() {
|
||||
mockPlsRepo.Data["file-partial"] = &model.Playlist{
|
||||
ID: "file-partial",
|
||||
Name: "Keep Me",
|
||||
OwnerID: "user-1",
|
||||
Path: "/music/p.m3u",
|
||||
Sync: true,
|
||||
}
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
err := repo.Update("file-partial", &model.Playlist{Sync: false}, "sync")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Keep Me"))
|
||||
Expect(mockPlsRepo.Last.Sync).To(BeFalse())
|
||||
})
|
||||
|
||||
It("renames the playlist when only name is sent", func() {
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
err := repo.Update("partial", &model.Playlist{Name: "Renamed"}, "name")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Renamed"))
|
||||
Expect(mockPlsRepo.Last.Comment).To(Equal("Original comment"))
|
||||
Expect(mockPlsRepo.Last.Public).To(BeFalse())
|
||||
})
|
||||
|
||||
It("clears the comment when an empty comment is sent explicitly", func() {
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
err := repo.Update("partial", &model.Playlist{Comment: ""}, "comment")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Comment).To(BeEmpty())
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Original Name"))
|
||||
})
|
||||
|
||||
It("updates rules-only on a smart playlist (Feishin-style edit)", func() {
|
||||
mockPlsRepo.Data["smart-partial"] = &model.Playlist{
|
||||
ID: "smart-partial",
|
||||
Name: "Smart Original",
|
||||
Comment: "smart comment",
|
||||
OwnerID: "user-1",
|
||||
Public: true,
|
||||
Rules: &criteria.Criteria{Expression: criteria.Is{"genre": "Rock"}},
|
||||
}
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
newRules := &criteria.Criteria{Expression: criteria.Is{"genre": "Jazz"}, Sort: "year DESC"}
|
||||
err := repo.Update("smart-partial", &model.Playlist{Rules: newRules}, "rules")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Rules).To(Equal(newRules))
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Smart Original"))
|
||||
Expect(mockPlsRepo.Last.Comment).To(Equal("smart comment"))
|
||||
Expect(mockPlsRepo.Last.Public).To(BeTrue())
|
||||
})
|
||||
|
||||
It("updates name and rules together (smart-playlist Edit form)", func() {
|
||||
mockPlsRepo.Data["smart-edit"] = &model.Playlist{
|
||||
ID: "smart-edit",
|
||||
Name: "Smart Original",
|
||||
Comment: "smart comment",
|
||||
OwnerID: "user-1",
|
||||
Rules: &criteria.Criteria{Expression: criteria.Is{"genre": "Rock"}},
|
||||
}
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
newRules := &criteria.Criteria{Expression: criteria.Is{"artist": "Miles Davis"}, Sort: "album"}
|
||||
err := repo.Update("smart-edit",
|
||||
&model.Playlist{Name: "Smart Renamed", Rules: newRules},
|
||||
"name", "rules")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Smart Renamed"))
|
||||
Expect(mockPlsRepo.Last.Rules).To(Equal(newRules))
|
||||
Expect(mockPlsRepo.Last.Comment).To(Equal("smart comment"))
|
||||
})
|
||||
|
||||
It("does not bump the saved rules on an idempotent rules-only PUT", func() {
|
||||
rules := &criteria.Criteria{Expression: criteria.Is{"genre": "Rock"}}
|
||||
mockPlsRepo.Data["smart-idempotent"] = &model.Playlist{
|
||||
ID: "smart-idempotent",
|
||||
Name: "Smart Idempotent",
|
||||
OwnerID: "user-1",
|
||||
Rules: rules,
|
||||
}
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
// Same rules sent back — rulesEqual should report no change and
|
||||
// the request should no-op (no Put call).
|
||||
sameRules := &criteria.Criteria{Expression: criteria.Is{"genre": "Rock"}}
|
||||
err := repo.Update("smart-idempotent", &model.Playlist{Rules: sameRules}, "rules")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last).To(BeNil()) // no Put happened
|
||||
})
|
||||
|
||||
It("preserves rules when only public is sent (smart playlist + bulk Make Public)", func() {
|
||||
rules := &criteria.Criteria{Expression: criteria.Is{"genre": "Rock"}}
|
||||
mockPlsRepo.Data["smart-public"] = &model.Playlist{
|
||||
ID: "smart-public",
|
||||
Name: "Smart Public",
|
||||
OwnerID: "user-1",
|
||||
Public: false,
|
||||
Rules: rules,
|
||||
}
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
err := repo.Update("smart-public", &model.Playlist{Public: true}, "public")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Public).To(BeTrue())
|
||||
Expect(mockPlsRepo.Last.Rules).To(Equal(rules))
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Smart Public"))
|
||||
})
|
||||
|
||||
It("does not treat a missing ownerId as an ownership transfer attempt", func() {
|
||||
// A non-admin user sending only {public:true} should not be blocked
|
||||
// just because OwnerID is the zero value in the deserialized entity.
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
err := repo.Update("partial", &model.Playlist{Public: true}, "public")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
})
|
||||
|
||||
It("matches cols case-insensitively (mirrors json decoder behavior)", func() {
|
||||
// Go's json decoder populates struct fields from case-variant keys
|
||||
// like {"Name":"x"}, but rest.Put's field-name extraction is
|
||||
// case-sensitive. sentFields normalizes both sides so a request
|
||||
// with {"Name":"Renamed"} is honored, not silently ignored.
|
||||
repo = ps.NewRepository(ctx).(rest.Persistable)
|
||||
err := repo.Update("partial", &model.Playlist{Name: "Renamed"}, "Name")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(mockPlsRepo.Last.Name).To(Equal("Renamed"))
|
||||
Expect(mockPlsRepo.Last.Comment).To(Equal("Original comment"))
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Delete", func() {
|
||||
|
||||
@ -98,7 +98,7 @@ func (r *shareRepositoryWrapper) Save(entity any) (string, error) {
|
||||
s.ExpiresAt = new(time.Now().Add(conf.Server.DefaultShareExpiration))
|
||||
}
|
||||
|
||||
firstId := strings.SplitN(s.ResourceIDs, ",", 2)[0]
|
||||
firstId, _, _ := strings.Cut(s.ResourceIDs, ",")
|
||||
v, err := model.GetEntityByID(r.ctx, r.ds, firstId)
|
||||
if err != nil {
|
||||
return "", err
|
||||
|
||||
135
core/stream/limiter.go
Normal file
135
core/stream/limiter.go
Normal file
@ -0,0 +1,135 @@
|
||||
package stream
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
)
|
||||
|
||||
// ErrTooManyTranscodes is returned by TranscodeLimiter.Acquire when the
|
||||
// configured concurrency cap has been reached. Callers should translate this
|
||||
// into an HTTP 429 response so well-behaved clients back off and retry.
|
||||
var ErrTooManyTranscodes = errors.New("too many concurrent transcodes")
|
||||
|
||||
// RetryAfterSeconds is the value returned in the HTTP Retry-After header when
|
||||
// a request is rejected with ErrTooManyTranscodes. Most transcodes finish well
|
||||
// within this window, so retrying after this delay typically succeeds.
|
||||
const RetryAfterSeconds = 5
|
||||
|
||||
// TranscodeLimiter gates the number of concurrent ffmpeg transcodes. It enforces
|
||||
// both a global cap (to protect the host from process exhaustion) and an optional
|
||||
// per-user cap (to keep one client from starving the others). Acquire never
|
||||
// blocks: it either reserves a slot or returns ErrTooManyTranscodes immediately.
|
||||
type TranscodeLimiter interface {
|
||||
// Acquire reserves a slot for the given user. On success it returns a release
|
||||
// function that must be called exactly once when the transcode is done.
|
||||
// Calling release more than once is safe and idempotent.
|
||||
Acquire(ctx context.Context, user string) (release func(), err error)
|
||||
|
||||
// Enabled reports whether the limiter actually enforces any cap. Callers
|
||||
// can use it to decide whether to bind ffmpeg's lifetime to the request
|
||||
// context so disconnects free slots quickly, rather than letting the
|
||||
// process drain to completion in the background.
|
||||
Enabled() bool
|
||||
}
|
||||
|
||||
// NewTranscodeLimiter returns a limiter enforcing the given caps. Each cap is
|
||||
// independent: a value of zero or less disables that cap. When both caps are
|
||||
// disabled the limiter is a no-op.
|
||||
func NewTranscodeLimiter(maxConcurrent, maxPerUser int) TranscodeLimiter {
|
||||
if maxConcurrent <= 0 && maxPerUser <= 0 {
|
||||
return noopLimiter{}
|
||||
}
|
||||
l := &transcodeLimiter{maxPerUser: maxPerUser}
|
||||
if maxConcurrent > 0 {
|
||||
l.global = make(chan struct{}, maxConcurrent)
|
||||
}
|
||||
if maxPerUser > 0 {
|
||||
l.perUser = make(map[string]int)
|
||||
}
|
||||
return l
|
||||
}
|
||||
|
||||
// releasingReadCloser wraps an io.ReadCloser so that closing it also releases
|
||||
// the limiter slot exactly once. release must be the function returned by
|
||||
// TranscodeLimiter.Acquire; its own idempotency makes double-Close safe too.
|
||||
type releasingReadCloser struct {
|
||||
io.ReadCloser
|
||||
release func()
|
||||
}
|
||||
|
||||
func (r *releasingReadCloser) Close() error {
|
||||
err := r.ReadCloser.Close()
|
||||
r.release()
|
||||
return err
|
||||
}
|
||||
|
||||
type noopLimiter struct{}
|
||||
|
||||
func (noopLimiter) Acquire(context.Context, string) (func(), error) {
|
||||
return func() {}, nil
|
||||
}
|
||||
|
||||
func (noopLimiter) Enabled() bool { return false }
|
||||
|
||||
type transcodeLimiter struct {
|
||||
maxPerUser int
|
||||
global chan struct{}
|
||||
|
||||
mu sync.Mutex
|
||||
perUser map[string]int
|
||||
}
|
||||
|
||||
func (*transcodeLimiter) Enabled() bool { return true }
|
||||
|
||||
func (l *transcodeLimiter) Acquire(_ context.Context, user string) (func(), error) {
|
||||
// Reserve a per-user slot first so a noisy user can't burn through
|
||||
// global slots only to be rejected later. An empty user key means
|
||||
// "anonymous" (e.g. public share viewers); we skip the per-user cap
|
||||
// entirely so unrelated anonymous clients do not share a bucket.
|
||||
perUserActive := l.maxPerUser > 0 && user != ""
|
||||
if perUserActive {
|
||||
l.mu.Lock()
|
||||
if l.perUser[user] >= l.maxPerUser {
|
||||
l.mu.Unlock()
|
||||
return nil, ErrTooManyTranscodes
|
||||
}
|
||||
l.perUser[user]++
|
||||
l.mu.Unlock()
|
||||
}
|
||||
|
||||
if l.global != nil {
|
||||
select {
|
||||
case l.global <- struct{}{}:
|
||||
default:
|
||||
if perUserActive {
|
||||
l.releasePerUser(user)
|
||||
}
|
||||
return nil, ErrTooManyTranscodes
|
||||
}
|
||||
}
|
||||
|
||||
var released atomic.Bool
|
||||
return func() {
|
||||
if !released.CompareAndSwap(false, true) {
|
||||
return
|
||||
}
|
||||
if l.global != nil {
|
||||
<-l.global
|
||||
}
|
||||
if perUserActive {
|
||||
l.releasePerUser(user)
|
||||
}
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (l *transcodeLimiter) releasePerUser(user string) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
l.perUser[user]--
|
||||
if l.perUser[user] <= 0 {
|
||||
delete(l.perUser, user)
|
||||
}
|
||||
}
|
||||
186
core/stream/limiter_test.go
Normal file
186
core/stream/limiter_test.go
Normal file
@ -0,0 +1,186 @@
|
||||
package stream_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
|
||||
"github.com/navidrome/navidrome/core/stream"
|
||||
"github.com/navidrome/navidrome/log"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
)
|
||||
|
||||
var _ = Describe("TranscodeLimiter", func() {
|
||||
ctx := log.NewContext(context.TODO())
|
||||
|
||||
Describe("Disabled (both caps <= 0)", func() {
|
||||
It("never blocks and never returns ErrTooManyTranscodes", func() {
|
||||
lim := stream.NewTranscodeLimiter(0, 0)
|
||||
for range 100 {
|
||||
rel, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(rel).ToNot(BeNil())
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Per-user cap only (no global cap)", func() {
|
||||
It("still enforces the per-user limit when MaxConcurrent is disabled", func() {
|
||||
lim := stream.NewTranscodeLimiter(0, 2)
|
||||
|
||||
rel1, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
rel2, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
_, err = lim.Acquire(ctx, "alice")
|
||||
Expect(errors.Is(err, stream.ErrTooManyTranscodes)).To(BeTrue())
|
||||
|
||||
// Other users have their own buckets.
|
||||
rel3, err := lim.Acquire(ctx, "bob")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
rel1()
|
||||
_, err = lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
rel2()
|
||||
rel3()
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Global cap", func() {
|
||||
It("rejects requests beyond MaxConcurrent with ErrTooManyTranscodes", func() {
|
||||
lim := stream.NewTranscodeLimiter(2, 0)
|
||||
|
||||
rel1, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
rel2, err := lim.Acquire(ctx, "bob")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
_, err = lim.Acquire(ctx, "carol")
|
||||
Expect(errors.Is(err, stream.ErrTooManyTranscodes)).To(BeTrue())
|
||||
|
||||
rel1()
|
||||
_, err = lim.Acquire(ctx, "carol")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
rel2()
|
||||
})
|
||||
|
||||
It("releases a slot only once even if release is called multiple times", func() {
|
||||
lim := stream.NewTranscodeLimiter(1, 0)
|
||||
|
||||
rel, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
rel()
|
||||
rel()
|
||||
rel()
|
||||
|
||||
// After releases, exactly one slot should be available.
|
||||
_, err = lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
_, err = lim.Acquire(ctx, "alice")
|
||||
Expect(errors.Is(err, stream.ErrTooManyTranscodes)).To(BeTrue())
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Per-user cap", func() {
|
||||
It("rejects a user beyond MaxConcurrentPerUser even if global slots remain", func() {
|
||||
lim := stream.NewTranscodeLimiter(10, 2)
|
||||
|
||||
rel1, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
rel2, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
_, err = lim.Acquire(ctx, "alice")
|
||||
Expect(errors.Is(err, stream.ErrTooManyTranscodes)).To(BeTrue())
|
||||
|
||||
// A different user is unaffected.
|
||||
rel3, err := lim.Acquire(ctx, "bob")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
rel1()
|
||||
_, err = lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
rel2()
|
||||
rel3()
|
||||
})
|
||||
|
||||
It("skips the per-user cap for anonymous users (empty key)", func() {
|
||||
// Anonymous requests (e.g. public share viewers) deliberately
|
||||
// bypass the per-user cap so unrelated anonymous clients are not
|
||||
// collapsed into a single shared bucket. The global cap remains
|
||||
// the only ceiling on anonymous traffic.
|
||||
lim := stream.NewTranscodeLimiter(10, 1)
|
||||
|
||||
rels := make([]func(), 0, 5)
|
||||
for range 5 {
|
||||
rel, err := lim.Acquire(ctx, "")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
rels = append(rels, rel)
|
||||
}
|
||||
for _, rel := range rels {
|
||||
rel()
|
||||
}
|
||||
})
|
||||
|
||||
It("still applies the global cap to anonymous users", func() {
|
||||
lim := stream.NewTranscodeLimiter(2, 1)
|
||||
|
||||
rel1, err := lim.Acquire(ctx, "")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
rel2, err := lim.Acquire(ctx, "")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
_, err = lim.Acquire(ctx, "")
|
||||
Expect(errors.Is(err, stream.ErrTooManyTranscodes)).To(BeTrue())
|
||||
|
||||
rel1()
|
||||
rel2()
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Concurrent safety", func() {
|
||||
It("survives parallel Acquire/release with consistent counts", func() {
|
||||
lim := stream.NewTranscodeLimiter(5, 0)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
var acquired int64
|
||||
var rejected int64
|
||||
var mu sync.Mutex
|
||||
|
||||
for i := range 50 {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
rel, err := lim.Acquire(ctx, "alice")
|
||||
mu.Lock()
|
||||
if err == nil {
|
||||
acquired++
|
||||
mu.Unlock()
|
||||
rel()
|
||||
} else {
|
||||
rejected++
|
||||
mu.Unlock()
|
||||
}
|
||||
_ = i
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
Expect(acquired + rejected).To(Equal(int64(50)))
|
||||
// After all releases, all 5 slots should be free again.
|
||||
for range 5 {
|
||||
_, err := lim.Acquire(ctx, "alice")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
}
|
||||
_, err := lim.Acquire(ctx, "alice")
|
||||
Expect(errors.Is(err, stream.ErrTooManyTranscodes)).To(BeTrue())
|
||||
})
|
||||
})
|
||||
})
|
||||
@ -2,6 +2,7 @@ package stream
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"mime"
|
||||
@ -28,13 +29,19 @@ type MediaStreamer interface {
|
||||
type TranscodingCache cache.FileCache
|
||||
|
||||
func NewMediaStreamer(ds model.DataStore, t ffmpeg.FFmpeg, cache TranscodingCache) MediaStreamer {
|
||||
return &mediaStreamer{ds: ds, transcoder: t, cache: cache}
|
||||
return &mediaStreamer{
|
||||
ds: ds,
|
||||
transcoder: t,
|
||||
cache: cache,
|
||||
limiter: NewTranscodeLimiter(conf.Server.Transcoding.MaxConcurrent, conf.Server.Transcoding.MaxConcurrentPerUser),
|
||||
}
|
||||
}
|
||||
|
||||
type mediaStreamer struct {
|
||||
ds model.DataStore
|
||||
transcoder ffmpeg.FFmpeg
|
||||
cache cache.FileCache
|
||||
limiter TranscodeLimiter
|
||||
}
|
||||
|
||||
type streamJob struct {
|
||||
@ -104,7 +111,12 @@ func (ms *mediaStreamer) NewStream(ctx context.Context, mf *model.MediaFile, req
|
||||
}
|
||||
r, err := ms.cache.Get(ctx, job)
|
||||
if err != nil {
|
||||
log.Error(ctx, "Error accessing transcoding cache", "id", mf.ID, err)
|
||||
// Rate-limit rejections are already logged at warn level by the
|
||||
// producer; treating them as cache failures here would both
|
||||
// double-log and mask actual cache problems.
|
||||
if !errors.Is(err, ErrTooManyTranscodes) {
|
||||
log.Error(ctx, "Error accessing transcoding cache", "id", mf.ID, err)
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
cached = r.Cached
|
||||
@ -217,15 +229,31 @@ func NewTranscodingCache() TranscodingCache {
|
||||
return nil, os.ErrInvalid
|
||||
}
|
||||
|
||||
// Choose the appropriate context based on EnableTranscodingCancellation configuration.
|
||||
// This is where we decide whether transcoding processes should be cancellable or not.
|
||||
release, err := job.ms.limiter.Acquire(ctx, limiterKey(ctx))
|
||||
if err != nil {
|
||||
log.Warn(ctx, "Refusing transcode: concurrent transcode limit reached",
|
||||
"id", job.mf.ID, "user", userName(ctx),
|
||||
"maxConcurrent", conf.Server.Transcoding.MaxConcurrent,
|
||||
"maxPerUser", conf.Server.Transcoding.MaxConcurrentPerUser)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Choose the context that drives the ffmpeg process.
|
||||
//
|
||||
// When the limiter is enabled, force the request context so a
|
||||
// client disconnect cancels ffmpeg and frees the slot promptly.
|
||||
// Otherwise a client could open many transcodes, disconnect
|
||||
// immediately, and still leave the configured cap's worth of
|
||||
// ffmpeg processes draining in the background — which is exactly
|
||||
// the DoS the limiter is meant to prevent.
|
||||
//
|
||||
// When the limiter is disabled, preserve the legacy behavior
|
||||
// governed by Transcoding.EnableCancellation so unchanged configs
|
||||
// keep their previous observable behavior.
|
||||
var transcodingCtx context.Context
|
||||
if conf.Server.EnableTranscodingCancellation {
|
||||
// Use the request context directly, allowing cancellation when client disconnects
|
||||
if job.ms.limiter.Enabled() || conf.Server.Transcoding.EnableCancellation {
|
||||
transcodingCtx = ctx
|
||||
} else {
|
||||
// Use background context with request values preserved.
|
||||
// This prevents cancellation but maintains request metadata (user, client, etc.)
|
||||
transcodingCtx = request.AddValues(context.Background(), ctx)
|
||||
}
|
||||
|
||||
@ -240,10 +268,14 @@ func NewTranscodingCache() TranscodingCache {
|
||||
Offset: job.offset,
|
||||
})
|
||||
if err != nil {
|
||||
release()
|
||||
log.Error(ctx, "Error starting transcoder", "id", job.mf.ID, err)
|
||||
return nil, os.ErrInvalid
|
||||
}
|
||||
return out, nil
|
||||
// Tie the slot to the ffmpeg process: copyAndClose calls Close
|
||||
// on this reader after io.Copy returns, which is exactly when
|
||||
// ffmpeg has exited (either EOF or context cancellation).
|
||||
return &releasingReadCloser{ReadCloser: out, release: release}, nil
|
||||
})
|
||||
}
|
||||
|
||||
@ -255,3 +287,16 @@ func userName(ctx context.Context) string {
|
||||
return user.UserName
|
||||
}
|
||||
}
|
||||
|
||||
// limiterKey returns the per-user bucket key used by the transcode limiter.
|
||||
// For anonymous requests (e.g. public shares) it returns the empty string,
|
||||
// which signals the limiter to skip the per-user cap entirely — otherwise
|
||||
// every anonymous viewer of a public share would collide on the same key
|
||||
// and starve each other within MaxConcurrentPerUser slots. The global cap
|
||||
// still applies and remains the protection against runaway anonymous load.
|
||||
func limiterKey(ctx context.Context) string {
|
||||
if user, ok := request.UserFrom(ctx); ok {
|
||||
return user.UserName
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
@ -2,6 +2,7 @@ package stream_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
|
||||
@ -10,6 +11,7 @@ import (
|
||||
"github.com/navidrome/navidrome/core/stream"
|
||||
"github.com/navidrome/navidrome/log"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/model/request"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
@ -61,6 +63,70 @@ var _ = Describe("MediaStreamer", func() {
|
||||
Expect(s.Seekable()).To(BeFalse())
|
||||
Expect(s.Duration()).To(Equal(float32(257.0)))
|
||||
})
|
||||
It("rejects transcode requests beyond MaxConcurrent with ErrTooManyTranscodes", func() {
|
||||
// Use an ffmpeg whose Read blocks indefinitely so the cache's
|
||||
// background copy can't drain the source and release the slot —
|
||||
// keeping the single transcode slot pinned for this test.
|
||||
pr, pw := io.Pipe()
|
||||
DeferCleanup(func() { _ = pw.Close() })
|
||||
blockingFFmpeg := tests.NewMockFFmpeg("")
|
||||
blockingFFmpeg.Reader = pr
|
||||
|
||||
conf.Server.Transcoding.MaxConcurrent = 1
|
||||
conf.Server.Transcoding.MaxConcurrentPerUser = 0
|
||||
tightCache := stream.NewTranscodingCache()
|
||||
Eventually(func() bool { return tightCache.Available(context.TODO()) }).Should(BeTrue())
|
||||
tightStreamer := stream.NewMediaStreamer(ds, blockingFFmpeg, tightCache)
|
||||
|
||||
userCtx := request.WithUsername(ctx, "alice")
|
||||
s1, err := tightStreamer.NewStream(userCtx, mf, stream.Request{Format: "mp3", BitRate: 64})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
defer s1.Close()
|
||||
|
||||
// Different cache key so it doesn't dedupe with the first request.
|
||||
_, err = tightStreamer.NewStream(userCtx, mf, stream.Request{Format: "mp3", BitRate: 96})
|
||||
Expect(errors.Is(err, stream.ErrTooManyTranscodes)).To(BeTrue())
|
||||
})
|
||||
|
||||
It("releases the slot once the stream is closed", func() {
|
||||
conf.Server.Transcoding.MaxConcurrent = 1
|
||||
conf.Server.Transcoding.MaxConcurrentPerUser = 0
|
||||
tightCache := stream.NewTranscodingCache()
|
||||
Eventually(func() bool { return tightCache.Available(context.TODO()) }).Should(BeTrue())
|
||||
tightStreamer := stream.NewMediaStreamer(ds, ffmpeg, tightCache)
|
||||
|
||||
userCtx := request.WithUsername(ctx, "alice")
|
||||
s1, err := tightStreamer.NewStream(userCtx, mf, stream.Request{Format: "mp3", BitRate: 64})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
_, _ = io.ReadAll(s1)
|
||||
_ = s1.Close()
|
||||
Eventually(func() bool { return ffmpeg.IsClosed() }, "3s").Should(BeTrue())
|
||||
|
||||
// Slot should now be free for a different transcode.
|
||||
s2, err := tightStreamer.NewStream(userCtx, mf, stream.Request{Format: "mp3", BitRate: 96})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
defer s2.Close()
|
||||
})
|
||||
|
||||
It("does not consume a slot for raw streams", func() {
|
||||
conf.Server.Transcoding.MaxConcurrent = 1
|
||||
conf.Server.Transcoding.MaxConcurrentPerUser = 0
|
||||
tightCache := stream.NewTranscodingCache()
|
||||
Eventually(func() bool { return tightCache.Available(context.TODO()) }).Should(BeTrue())
|
||||
tightStreamer := stream.NewMediaStreamer(ds, ffmpeg, tightCache)
|
||||
|
||||
userCtx := request.WithUsername(ctx, "alice")
|
||||
// First, saturate the single transcode slot.
|
||||
s1, err := tightStreamer.NewStream(userCtx, mf, stream.Request{Format: "mp3", BitRate: 64})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
defer s1.Close()
|
||||
|
||||
// Raw stream must still succeed.
|
||||
s2, err := tightStreamer.NewStream(userCtx, mf, stream.Request{Format: "raw"})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
defer s2.Close()
|
||||
})
|
||||
|
||||
It("returns a seekable stream if the file is complete in the cache", func() {
|
||||
s, err := streamer.NewStream(ctx, mf, stream.Request{Format: "mp3", BitRate: 32})
|
||||
Expect(err).To(BeNil())
|
||||
|
||||
18
go.mod
18
go.mod
@ -20,7 +20,7 @@ require (
|
||||
github.com/extism/go-sdk v1.7.1
|
||||
github.com/fatih/structs v1.1.0
|
||||
github.com/gen2brain/webp v0.5.5
|
||||
github.com/go-chi/chi/v5 v5.2.5
|
||||
github.com/go-chi/chi/v5 v5.3.0
|
||||
github.com/go-chi/cors v1.2.2
|
||||
github.com/go-chi/httprate v0.15.0
|
||||
github.com/go-chi/jwtauth/v5 v5.4.0
|
||||
@ -39,8 +39,8 @@ require (
|
||||
github.com/mattn/go-sqlite3 v1.14.44
|
||||
github.com/microcosm-cc/bluemonday v1.0.27
|
||||
github.com/mileusna/useragent v1.3.5
|
||||
github.com/onsi/ginkgo/v2 v2.28.3
|
||||
github.com/onsi/gomega v1.40.0
|
||||
github.com/onsi/ginkgo/v2 v2.29.0
|
||||
github.com/onsi/gomega v1.41.0
|
||||
github.com/pelletier/go-toml/v2 v2.3.1
|
||||
github.com/pmezard/go-difflib v1.0.0
|
||||
github.com/pocketbase/dbx v1.12.0
|
||||
@ -59,10 +59,10 @@ require (
|
||||
github.com/xrash/smetrics v0.0.0-20250705151800-55b8f293f342
|
||||
go.senan.xyz/taglib v0.11.1
|
||||
go.uber.org/goleak v1.3.0
|
||||
golang.org/x/image v0.40.0
|
||||
golang.org/x/net v0.54.0
|
||||
golang.org/x/image v0.41.0
|
||||
golang.org/x/net v0.55.0
|
||||
golang.org/x/sync v0.20.0
|
||||
golang.org/x/sys v0.44.0
|
||||
golang.org/x/sys v0.45.0
|
||||
golang.org/x/term v0.43.0
|
||||
golang.org/x/text v0.37.0
|
||||
golang.org/x/time v0.15.0
|
||||
@ -81,7 +81,7 @@ require (
|
||||
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
|
||||
github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.1 // indirect
|
||||
github.com/dylibso/observe-sdk/go v0.0.0-20240828172851-9145d8ad07e1 // indirect
|
||||
github.com/ebitengine/purego v0.10.0 // indirect
|
||||
github.com/ebitengine/purego v0.10.1 // indirect
|
||||
github.com/fsnotify/fsnotify v1.10.1 // indirect
|
||||
github.com/go-logr/logr v1.4.3 // indirect
|
||||
github.com/go-task/slim-sprig/v3 v3.0.0 // indirect
|
||||
@ -115,7 +115,7 @@ require (
|
||||
github.com/prometheus/client_model v0.6.2 // indirect
|
||||
github.com/prometheus/common v0.67.5 // indirect
|
||||
github.com/prometheus/procfs v0.20.1 // indirect
|
||||
github.com/rogpeppe/go-internal v1.14.1 // indirect
|
||||
github.com/rogpeppe/go-internal v1.15.0 // indirect
|
||||
github.com/sagikazarmark/locafero v0.12.0 // indirect
|
||||
github.com/sanity-io/litter v1.5.8 // indirect
|
||||
github.com/segmentio/asm v1.2.1 // indirect
|
||||
@ -133,7 +133,7 @@ require (
|
||||
go.uber.org/multierr v1.11.0 // indirect
|
||||
go.yaml.in/yaml/v2 v2.4.3 // indirect
|
||||
go.yaml.in/yaml/v3 v3.0.4 // indirect
|
||||
golang.org/x/crypto v0.51.0 // indirect
|
||||
golang.org/x/crypto v0.52.0 // indirect
|
||||
golang.org/x/mod v0.36.0 // indirect
|
||||
golang.org/x/telemetry v0.0.0-20260508192327-42602be52be6 // indirect
|
||||
golang.org/x/tools v0.45.0 // indirect
|
||||
|
||||
36
go.sum
36
go.sum
@ -54,8 +54,8 @@ github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkp
|
||||
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
|
||||
github.com/dylibso/observe-sdk/go v0.0.0-20240828172851-9145d8ad07e1 h1:idfl8M8rPW93NehFw5H1qqH8yG158t5POr+LX9avbJY=
|
||||
github.com/dylibso/observe-sdk/go v0.0.0-20240828172851-9145d8ad07e1/go.mod h1:C8DzXehI4zAbrdlbtOByKX6pfivJTBiV9Jjqv56Yd9Q=
|
||||
github.com/ebitengine/purego v0.10.0 h1:QIw4xfpWT6GWTzaW5XEKy3HXoqrJGx1ijYHzTF0/ISU=
|
||||
github.com/ebitengine/purego v0.10.0/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ=
|
||||
github.com/ebitengine/purego v0.10.1 h1:dewVBCBT2GaMu1SrNTYxQhgQBethzfhiwvZiLGP/qyY=
|
||||
github.com/ebitengine/purego v0.10.1/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ=
|
||||
github.com/extism/go-sdk v1.7.1 h1:lWJos6uY+tRFdlIHR+SJjwFDApY7OypS/2nMhiVQ9Sw=
|
||||
github.com/extism/go-sdk v1.7.1/go.mod h1:IT+Xdg5AZM9hVtpFUA+uZCJMge/hbvshl8bwzLtFyKA=
|
||||
github.com/fatih/structs v1.1.0 h1:Q7juDM0QtcnhCpeyLGQKyg4TOIghuNXrkL32pHAUMxo=
|
||||
@ -73,8 +73,8 @@ github.com/gkampitakis/go-diff v1.3.2 h1:Qyn0J9XJSDTgnsgHRdz9Zp24RaJeKMUHg2+PDZZ
|
||||
github.com/gkampitakis/go-diff v1.3.2/go.mod h1:LLgOrpqleQe26cte8s36HTWcTmMEur6OPYerdAAS9tk=
|
||||
github.com/gkampitakis/go-snaps v0.5.15 h1:amyJrvM1D33cPHwVrjo9jQxX8g/7E2wYdZ+01KS3zGE=
|
||||
github.com/gkampitakis/go-snaps v0.5.15/go.mod h1:HNpx/9GoKisdhw9AFOBT1N7DBs9DiHo/hGheFGBZ+mc=
|
||||
github.com/go-chi/chi/v5 v5.2.5 h1:Eg4myHZBjyvJmAFjFvWgrqDTXFyOzjj7YIm3L3mu6Ug=
|
||||
github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0=
|
||||
github.com/go-chi/chi/v5 v5.3.0 h1:halUjDxhshgXHMrao5bB8eNBXo/rnzwr8m5m36glehM=
|
||||
github.com/go-chi/chi/v5 v5.3.0/go.mod h1:R+tYY2hNuVUUjxoPtqUdgBqevM9s9njzkTLutVsOCto=
|
||||
github.com/go-chi/cors v1.2.2 h1:Jmey33TE+b+rB7fT8MUy1u0I4L+NARQlK6LhzKPSyQE=
|
||||
github.com/go-chi/cors v1.2.2/go.mod h1:sSbTewc+6wYHBBCW7ytsFSn836hqM7JxpglAy2Vzc58=
|
||||
github.com/go-chi/httprate v0.15.0 h1:j54xcWV9KGmPf/X4H32/aTH+wBlrvxL7P+SdnRqxh5g=
|
||||
@ -193,10 +193,10 @@ github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOF
|
||||
github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls=
|
||||
github.com/ogier/pflag v0.0.1 h1:RW6JSWSu/RkSatfcLtogGfFgpim5p7ARQ10ECk5O750=
|
||||
github.com/ogier/pflag v0.0.1/go.mod h1:zkFki7tvTa0tafRvTBIZTvzYyAu6kQhPZFnshFFPE+g=
|
||||
github.com/onsi/ginkgo/v2 v2.28.3 h1:4JvMdwtFU0imd8fHx25OJXoDMRexnf8v5NHKYSTTji4=
|
||||
github.com/onsi/ginkgo/v2 v2.28.3/go.mod h1:+aXOY+vzZ5mu2iI2HpTZUPmM//oQfsNFX6gU9kNcA44=
|
||||
github.com/onsi/gomega v1.40.0 h1:Vtol0e1MghCD2ZVIilPDIg44XSL9l2QAn8ZNaljWcJc=
|
||||
github.com/onsi/gomega v1.40.0/go.mod h1:M/Uqpu/8qTjtzCLUA2zJHX9Iilrau25x1PdoSRbWh5A=
|
||||
github.com/onsi/ginkgo/v2 v2.29.0 h1:rfh+ZFjgJhYWRoIqVf3Uwx/W20yLrcrE2h2GmYVRaag=
|
||||
github.com/onsi/ginkgo/v2 v2.29.0/go.mod h1:+aXOY+vzZ5mu2iI2HpTZUPmM//oQfsNFX6gU9kNcA44=
|
||||
github.com/onsi/gomega v1.41.0 h1:OwKp4pXNgVxf6sCplzYo794OFNuoL2q2SBMU5NSWOjA=
|
||||
github.com/onsi/gomega v1.41.0/go.mod h1:M/Uqpu/8qTjtzCLUA2zJHX9Iilrau25x1PdoSRbWh5A=
|
||||
github.com/pelletier/go-toml/v2 v2.3.1 h1:MYEvvGnQjeNkRF1qUuGolNtNExTDwct51yp7olPtrEc=
|
||||
github.com/pelletier/go-toml/v2 v2.3.1/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY=
|
||||
github.com/pkg/diff v0.0.0-20210226163009-20ebb0f2a09e/go.mod h1:pJLUxLENpZxwdsKMEsNbx1VGcRFpLqf3715MtcvvzbA=
|
||||
@ -224,8 +224,8 @@ github.com/rjeczalik/notify v0.9.3/go.mod h1:gF3zSOrafR9DQEWSE8TjfI9NkooDxbyT4Ug
|
||||
github.com/robfig/cron/v3 v3.0.1 h1:WdRxkvbJztn8LMz/QEvLN5sBU+xKpSqwwUO1Pjr4qDs=
|
||||
github.com/robfig/cron/v3 v3.0.1/go.mod h1:eQICP3HwyT7UooqI/z+Ov+PtYAWygg1TEWWzGIFLtro=
|
||||
github.com/rogpeppe/go-internal v1.9.0/go.mod h1:WtVeX8xhTBvf0smdhujwtBcq4Qrzq/fJaraNFVN+nFs=
|
||||
github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ=
|
||||
github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc=
|
||||
github.com/rogpeppe/go-internal v1.15.0 h1:D0RCU5rMAp+SpgkiNdrjfJ+LX4J1M32V2NeCY7EJ6hc=
|
||||
github.com/rogpeppe/go-internal v1.15.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs=
|
||||
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
|
||||
github.com/sabhiram/go-gitignore v0.0.0-20210923224102-525f6e181f06 h1:OkMGxebDjyw0ULyrTYWeN0UNCCkmCWfjPnIA2W6oviI=
|
||||
github.com/sabhiram/go-gitignore v0.0.0-20210923224102-525f6e181f06/go.mod h1:+ePHsJ1keEjQtpvf9HHw0f4ZeJ0TLRsxhunSI2hYJSs=
|
||||
@ -316,10 +316,10 @@ golang.org/x/crypto v0.13.0/go.mod h1:y6Z2r+Rw4iayiXXAIxJIDAJ1zMW4yaTpebo8fPOliY
|
||||
golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU=
|
||||
golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8=
|
||||
golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk=
|
||||
golang.org/x/crypto v0.51.0 h1:IBPXwPfKxY7cWQZ38ZCIRPI50YLeevDLlLnyC5wRGTI=
|
||||
golang.org/x/crypto v0.51.0/go.mod h1:8AdwkbraGNABw2kOX6YFPs3WM22XqI4EXEd8g+x7Oc8=
|
||||
golang.org/x/image v0.40.0 h1:Tw4GyDXMo+daZN1znreBRC3VayR1aLFUyUEOLUdW1a8=
|
||||
golang.org/x/image v0.40.0/go.mod h1:uIc348UZMSvS5Z65CVZ7iDPaNobNFEPeJ4kbqTOszmA=
|
||||
golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988=
|
||||
golang.org/x/crypto v0.52.0/go.mod h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc=
|
||||
golang.org/x/image v0.41.0 h1:8wS72eGJMJaBxK6okTzd4WaXumUlTVlb753MlsSvTCo=
|
||||
golang.org/x/image v0.41.0/go.mod h1:uIc348UZMSvS5Z65CVZ7iDPaNobNFEPeJ4kbqTOszmA=
|
||||
golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4=
|
||||
golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
|
||||
golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
|
||||
@ -338,8 +338,8 @@ golang.org/x/net v0.15.0/go.mod h1:idbUs1IY1+zTqbi8yxTbhexhEEk5ur9LInksu6HrEpk=
|
||||
golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44=
|
||||
golang.org/x/net v0.25.0/go.mod h1:JkAGAh7GEvH74S6FOH42FLoXpXbE/aqXSrIQjXgsiwM=
|
||||
golang.org/x/net v0.33.0/go.mod h1:HXLR5J+9DxmrqMwG9qjGCxZ+zKXxBru04zlTvWlWuN4=
|
||||
golang.org/x/net v0.54.0 h1:2zJIZAxAHV/OHCDTCOHAYehQzLfSXuf/5SoL/Dv6w/w=
|
||||
golang.org/x/net v0.54.0/go.mod h1:Sj4oj8jK6XmHpBZU/zWHw3BV3abl4Kvi+Ut7cQcY+cQ=
|
||||
golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8=
|
||||
golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww=
|
||||
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
@ -364,8 +364,8 @@ golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.28.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.44.0 h1:ildZl3J4uzeKP07r2F++Op7E9B29JRUy+a27EibtBTQ=
|
||||
golang.org/x/sys v0.44.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
|
||||
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE=
|
||||
golang.org/x/telemetry v0.0.0-20260508192327-42602be52be6 h1:HjU6IWBiAgRIdAJ9/y1rwCn+UELEmwV+VsTLzj/W4sE=
|
||||
golang.org/x/telemetry v0.0.0-20260508192327-42602be52be6/go.mod h1:Eqhaxk/wZsWEH8CRxLwj6xzEJbz7k1EFGqx7nyCoabE=
|
||||
|
||||
@ -36,6 +36,6 @@ func (f *journalFormatter) Format(entry *logrus.Entry) ([]byte, error) {
|
||||
if !ok {
|
||||
priority = 6 // default to info for unknown levels
|
||||
}
|
||||
prefix := []byte(fmt.Sprintf("<%d>", priority))
|
||||
prefix := fmt.Appendf(nil, "<%d>", priority)
|
||||
return append(prefix, formatted...), nil
|
||||
}
|
||||
|
||||
@ -684,6 +684,26 @@ var _ = Describe("Participants", func() {
|
||||
Expect(composers[2].Name).To(Equal("The Album Artist"))
|
||||
})
|
||||
})
|
||||
|
||||
// Sibling fix to https://github.com/navidrome/navidrome/issues/5065: when
|
||||
// multiple frames map to the same role tag (e.g. TIPL producer entries),
|
||||
// the configured split separator must still apply to each value.
|
||||
When("the tag has multiple values", func() {
|
||||
It("should split each value individually", func() {
|
||||
mf = toMediaFile(model.RawTags{
|
||||
"COMPOSER": {"John Doe/Jane Doe", "Someone Else"},
|
||||
})
|
||||
|
||||
participants := mf.Participants
|
||||
Expect(participants).To(HaveKeyWithValue(model.RoleComposer, HaveLen(3)))
|
||||
composers := participants[model.RoleComposer]
|
||||
Expect(composers).To(ConsistOf(
|
||||
HaveField("Name", "John Doe"),
|
||||
HaveField("Name", "Jane Doe"),
|
||||
HaveField("Name", "Someone Else"),
|
||||
))
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
Describe("MBID tags", func() {
|
||||
|
||||
@ -129,6 +129,21 @@ var _ = Describe("Metadata", func() {
|
||||
|
||||
Expect(md.Strings(model.TagGenre)).To(Equal([]string{"Rock", "Pop", "Punk"}))
|
||||
})
|
||||
|
||||
// Regression test for https://github.com/navidrome/navidrome/issues/5065
|
||||
//
|
||||
// MP3s with both an ID3v2 TMOO frame and a TXXX:MOOD frame are surfaced by
|
||||
// TagLib's PropertyMap as a single "mood" key with multiple values. The split
|
||||
// configuration must still apply to each value individually.
|
||||
It("should split values from multiple frames mapping to the same tag", func() {
|
||||
props.Tags = model.RawTags{
|
||||
// Same shape as the bug report: two frames, comma-separated content.
|
||||
"mood": {"Love, Emotional, Ballad", "Love; Emotional; Ballad"},
|
||||
}
|
||||
md = metadata.New(filePath, props)
|
||||
|
||||
Expect(md.Strings(model.TagMood)).To(ConsistOf("Love", "Emotional", "Ballad"))
|
||||
})
|
||||
})
|
||||
|
||||
DescribeTable("Date",
|
||||
|
||||
@ -34,23 +34,25 @@ type TagConf struct {
|
||||
SplitRx *regexp.Regexp `yaml:"-"`
|
||||
}
|
||||
|
||||
// SplitTagValue splits a tag value by the split separators, but only if it has a single value.
|
||||
// SplitTagValue splits tag values by the configured split separators.
|
||||
// Each value in the input slice is individually split and trimmed.
|
||||
func (c TagConf) SplitTagValue(values []string) []string {
|
||||
// If there's not exactly one value or no separators, return early.
|
||||
if len(values) != 1 || c.SplitRx == nil {
|
||||
if c.SplitRx == nil || len(values) == 0 {
|
||||
return values
|
||||
}
|
||||
tag := values[0]
|
||||
|
||||
// Replace all occurrences of any separator with the zero-width space.
|
||||
tag = c.SplitRx.ReplaceAllString(tag, consts.Zwsp)
|
||||
var result []string
|
||||
for _, tag := range values {
|
||||
// Replace all occurrences of any separator with the zero-width space.
|
||||
tag = c.SplitRx.ReplaceAllString(tag, consts.Zwsp)
|
||||
|
||||
// Split by the zero-width space and trim each substring.
|
||||
parts := strings.Split(tag, consts.Zwsp)
|
||||
for i, part := range parts {
|
||||
parts[i] = strings.TrimSpace(part)
|
||||
// Split by the zero-width space and trim each substring.
|
||||
parts := strings.SplitSeq(tag, consts.Zwsp)
|
||||
for part := range parts {
|
||||
result = append(result, strings.TrimSpace(part))
|
||||
}
|
||||
}
|
||||
return parts
|
||||
return result
|
||||
}
|
||||
|
||||
type TagType string
|
||||
|
||||
64
model/tag_mappings_test.go
Normal file
64
model/tag_mappings_test.go
Normal file
@ -0,0 +1,64 @@
|
||||
package model
|
||||
|
||||
import (
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
)
|
||||
|
||||
var _ = Describe("TagConf", func() {
|
||||
Describe("SplitTagValue", func() {
|
||||
var conf TagConf
|
||||
|
||||
BeforeEach(func() {
|
||||
conf = TagConf{Split: []string{";", "/", ","}}
|
||||
conf.SplitRx = compileSplitRegex("test", conf.Split)
|
||||
})
|
||||
|
||||
It("splits a single value on configured separators", func() {
|
||||
Expect(conf.SplitTagValue([]string{"Rock/Pop;Punk"})).To(Equal([]string{"Rock", "Pop", "Punk"}))
|
||||
})
|
||||
|
||||
It("trims whitespace around split values", func() {
|
||||
Expect(conf.SplitTagValue([]string{"Love, Emotional, Ballad"})).To(Equal([]string{"Love", "Emotional", "Ballad"}))
|
||||
})
|
||||
|
||||
// Regression test for https://github.com/navidrome/navidrome/issues/5065
|
||||
//
|
||||
// When multiple ID3v2 frames map to the same logical tag (e.g. TMOO + TXXX:MOOD),
|
||||
// TagLib's PropertyMap merges them into a slice with several entries. Previously
|
||||
// SplitTagValue had a `len(values) != 1` guard that skipped splitting in this case.
|
||||
It("splits each value individually when given multiple inputs", func() {
|
||||
input := []string{"Love, Emotional, Ballad", "Love; Emotional; Ballad"}
|
||||
Expect(conf.SplitTagValue(input)).To(Equal([]string{
|
||||
"Love", "Emotional", "Ballad",
|
||||
"Love", "Emotional", "Ballad",
|
||||
}))
|
||||
})
|
||||
|
||||
It("matches separators case-insensitively when the split pattern allows", func() {
|
||||
c := TagConf{Split: []string{" AND "}}
|
||||
c.SplitRx = compileSplitRegex("test", c.Split)
|
||||
Expect(c.SplitTagValue([]string{"foo and bar AND baz"})).To(Equal([]string{"foo", "bar", "baz"}))
|
||||
})
|
||||
|
||||
It("returns values unchanged when no separators are configured", func() {
|
||||
c := TagConf{}
|
||||
Expect(c.SplitTagValue([]string{"Foo, Bar"})).To(Equal([]string{"Foo, Bar"}))
|
||||
Expect(c.SplitTagValue([]string{"a", "b"})).To(Equal([]string{"a", "b"}))
|
||||
})
|
||||
|
||||
It("returns an empty slice for empty input", func() {
|
||||
Expect(conf.SplitTagValue([]string{})).To(BeEmpty())
|
||||
})
|
||||
|
||||
It("handles a value with no separator as a single-element result", func() {
|
||||
Expect(conf.SplitTagValue([]string{"JustOneMood"})).To(Equal([]string{"JustOneMood"}))
|
||||
})
|
||||
|
||||
It("produces empty strings when separators are adjacent (dedup happens downstream)", func() {
|
||||
// SplitTagValue itself does not filter empties; that is the job of
|
||||
// filterDuplicatedOrEmptyValues in the metadata pipeline.
|
||||
Expect(conf.SplitTagValue([]string{"Rock//Pop"})).To(Equal([]string{"Rock", "", "Pop"}))
|
||||
})
|
||||
})
|
||||
})
|
||||
@ -66,7 +66,7 @@ func normalizeForFTS(values ...string) string {
|
||||
result = append(result, variant)
|
||||
}
|
||||
for _, v := range values {
|
||||
for _, word := range strings.Fields(v) {
|
||||
for word := range strings.FieldsSeq(v) {
|
||||
transliterated := sanitize.Accents(word)
|
||||
// Concatenated ASCII form: R.E.M. → REM, AC/DC → ACDC, St-Étienne → StEtienne.
|
||||
add(word, fts5PunctStrip.ReplaceAllString(transliterated, ""))
|
||||
@ -279,9 +279,9 @@ type ftsSearch struct {
|
||||
}
|
||||
|
||||
// ToSql returns a single-query fallback for the REST filter path (no two-phase split).
|
||||
func (s *ftsSearch) ToSql() (string, []interface{}, error) {
|
||||
func (s *ftsSearch) ToSql() (string, []any, error) {
|
||||
sql := s.tableName + ".rowid IN (SELECT rowid FROM " + s.ftsTable + " WHERE " + s.ftsTable + " MATCH ?)"
|
||||
return sql, []interface{}{s.matchExpr}, nil
|
||||
return sql, []any{s.matchExpr}, nil
|
||||
}
|
||||
|
||||
// execute runs a two-phase FTS5 search:
|
||||
@ -373,8 +373,8 @@ func ftsQueryDegraded(original, ftsQuery string) bool {
|
||||
// Check if all effective FTS tokens are very short (≤2 chars).
|
||||
// Short tokens with prefix matching are too broad when special chars were stripped.
|
||||
// For quoted phrases, extract the content and check the tokens inside.
|
||||
tokens := strings.Fields(ftsQuery)
|
||||
for _, t := range tokens {
|
||||
tokens := strings.FieldsSeq(ftsQuery)
|
||||
for t := range tokens {
|
||||
t = strings.TrimSuffix(t, "*")
|
||||
// Skip internal phrase placeholders
|
||||
if strings.HasPrefix(t, "\x00") {
|
||||
@ -390,7 +390,7 @@ func ftsQueryDegraded(original, ftsQuery string) bool {
|
||||
// Extract content between quotes
|
||||
inner := strings.Trim(t, `"`)
|
||||
innerAlpha := fts5PunctStrip.ReplaceAllString(inner, " ")
|
||||
for _, it := range strings.Fields(innerAlpha) {
|
||||
for it := range strings.FieldsSeq(innerAlpha) {
|
||||
if len(it) > 2 {
|
||||
return false
|
||||
}
|
||||
|
||||
@ -16,7 +16,7 @@ type likeSearch struct {
|
||||
filter Sqlizer
|
||||
}
|
||||
|
||||
func (s *likeSearch) ToSql() (string, []interface{}, error) {
|
||||
func (s *likeSearch) ToSql() (string, []any, error) {
|
||||
return s.filter.ToSql()
|
||||
}
|
||||
|
||||
|
||||
@ -6,6 +6,7 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"maps"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
@ -540,9 +541,7 @@ func (s *taskQueueServiceImpl) cleanupLoop() {
|
||||
func (s *taskQueueServiceImpl) runCleanup() {
|
||||
s.mu.Lock()
|
||||
queues := make(map[string]*queueState, len(s.queues))
|
||||
for k, v := range s.queues {
|
||||
queues[k] = v
|
||||
}
|
||||
maps.Copy(queues, s.queues)
|
||||
s.mu.Unlock()
|
||||
|
||||
now := time.Now().UnixMilli()
|
||||
|
||||
@ -367,8 +367,8 @@ var _ = Describe("TaskQueueService", func() {
|
||||
|
||||
// Enqueue several more tasks — they stay pending since the worker is busy
|
||||
var pendingIDs []string
|
||||
for i := 0; i < 3; i++ {
|
||||
taskID, err := service.Enqueue(ctx, "clear-test", []byte(fmt.Sprintf("task-%d", i)))
|
||||
for i := range 3 {
|
||||
taskID, err := service.Enqueue(ctx, "clear-test", fmt.Appendf(nil, "task-%d", i))
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
pendingIDs = append(pendingIDs, taskID)
|
||||
}
|
||||
@ -674,8 +674,8 @@ var _ = Describe("TaskQueueService", func() {
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Enqueue 5 tasks
|
||||
for i := 0; i < 5; i++ {
|
||||
_, err := service.Enqueue(ctx, "delay-concurrent", []byte(fmt.Sprintf("task-%d", i)))
|
||||
for i := range 5 {
|
||||
_, err := service.Enqueue(ctx, "delay-concurrent", fmt.Appendf(nil, "task-%d", i))
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
}
|
||||
|
||||
@ -1112,7 +1112,7 @@ var _ = Describe("TaskQueueService Integration", Ordered, func() {
|
||||
// the second will be dequeued but block on the rate limiter (status=running),
|
||||
// the rest will stay pending.
|
||||
var taskIDs []string
|
||||
for i := 0; i < 5; i++ {
|
||||
for range 5 {
|
||||
output, err := callTestTaskQueue(ctx, testTaskQueueInput{
|
||||
Operation: "enqueue",
|
||||
QueueName: "test-cancel",
|
||||
@ -1186,11 +1186,11 @@ var _ = Describe("TaskQueueService Integration", Ordered, func() {
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Enqueue several tasks
|
||||
for i := 0; i < 4; i++ {
|
||||
for i := range 4 {
|
||||
_, err := callTestTaskQueue(ctx, testTaskQueueInput{
|
||||
Operation: "enqueue",
|
||||
QueueName: "test-clear",
|
||||
Payload: []byte(fmt.Sprintf("task-%d", i)),
|
||||
Payload: fmt.Appendf(nil, "task-%d", i),
|
||||
})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
}
|
||||
|
||||
@ -185,7 +185,7 @@ var _ = Describe("ParseCrontab", func() {
|
||||
// findSetBit returns the lowest bit position set in v, ignoring the starBit (bit 63).
|
||||
func findSetBit(v uint64) int {
|
||||
v &^= 1 << 63 // clear starBit
|
||||
for i := 0; i < 63; i++ {
|
||||
for i := range 63 {
|
||||
if v&(1<<uint(i)) != 0 {
|
||||
return i
|
||||
}
|
||||
|
||||
@ -7,7 +7,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/navidrome/navidrome/core/auth"
|
||||
"github.com/navidrome/navidrome/core/stream"
|
||||
streampkg "github.com/navidrome/navidrome/core/stream"
|
||||
"github.com/navidrome/navidrome/log"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
. "github.com/navidrome/navidrome/utils/gg"
|
||||
@ -48,10 +48,15 @@ func (pub *Router) handleStream(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
stream, err := pub.streamer.NewStream(ctx, mf, stream.Request{
|
||||
stream, err := pub.streamer.NewStream(ctx, mf, streampkg.Request{
|
||||
Format: info.format, BitRate: info.bitrate,
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, streampkg.ErrTooManyTranscodes) {
|
||||
w.Header().Set("Retry-After", strconv.Itoa(streampkg.RetryAfterSeconds))
|
||||
http.Error(w, "too many concurrent transcodes, please retry shortly", http.StatusTooManyRequests)
|
||||
return
|
||||
}
|
||||
log.Error(ctx, "Error starting shared stream", err)
|
||||
http.Error(w, "invalid request", http.StatusInternalServerError)
|
||||
return
|
||||
|
||||
@ -7,6 +7,7 @@ import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"regexp"
|
||||
"strconv"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
@ -304,6 +305,8 @@ func mapToSubsonicError(err error) subError {
|
||||
err = newError(responses.ErrorDataNotFound, "data not found")
|
||||
case errors.Is(err, model.ErrNotAuthorized):
|
||||
err = newError(responses.ErrorAuthorizationFail)
|
||||
case errors.Is(err, stream.ErrTooManyTranscodes):
|
||||
err = newError(responses.ErrorGeneric, "too many concurrent transcodes, please retry shortly")
|
||||
default:
|
||||
err = newError(responses.ErrorGeneric, fmt.Sprintf("Internal Server Error: %s", err))
|
||||
}
|
||||
@ -313,15 +316,31 @@ func mapToSubsonicError(err error) subError {
|
||||
}
|
||||
|
||||
func sendError(w http.ResponseWriter, r *http.Request, err error) {
|
||||
if errors.Is(err, stream.ErrTooManyTranscodes) {
|
||||
w.Header().Set("Retry-After", strconv.Itoa(stream.RetryAfterSeconds))
|
||||
sendResponseWithStatus(w, r, errorResponse(err), http.StatusTooManyRequests)
|
||||
return
|
||||
}
|
||||
sendResponse(w, r, errorResponse(err))
|
||||
}
|
||||
|
||||
func errorResponse(err error) *responses.Subsonic {
|
||||
subErr := mapToSubsonicError(err)
|
||||
response := newResponse()
|
||||
response.Status = responses.StatusFailed
|
||||
response.Error = &responses.Error{Code: subErr.code, Message: subErr.Error()}
|
||||
|
||||
sendResponse(w, r, response)
|
||||
return response
|
||||
}
|
||||
|
||||
func sendResponse(w http.ResponseWriter, r *http.Request, payload *responses.Subsonic) {
|
||||
sendResponseWithStatus(w, r, payload, 0)
|
||||
}
|
||||
|
||||
// sendResponseWithStatus writes the response body in the format requested by
|
||||
// the client. When status is non-zero, WriteHeader is called with that code
|
||||
// before the body is written; callers that need to set additional headers
|
||||
// (e.g. Retry-After) must set them before calling.
|
||||
func sendResponseWithStatus(w http.ResponseWriter, r *http.Request, payload *responses.Subsonic, status int) {
|
||||
p := req.Params(r)
|
||||
f, _ := p.String("f")
|
||||
var response []byte
|
||||
@ -356,6 +375,9 @@ func sendResponse(w http.ResponseWriter, r *http.Request, payload *responses.Sub
|
||||
sendError(w, r, err)
|
||||
return
|
||||
}
|
||||
if status != 0 {
|
||||
w.WriteHeader(status)
|
||||
}
|
||||
|
||||
if payload.Status == responses.StatusOK {
|
||||
if log.IsGreaterOrEqualTo(log.LevelTrace) {
|
||||
|
||||
@ -1,17 +1,19 @@
|
||||
package subsonic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"encoding/xml"
|
||||
"fmt"
|
||||
"math"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
|
||||
"github.com/navidrome/navidrome/core/stream"
|
||||
"github.com/navidrome/navidrome/server/subsonic/responses"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
||||
var _ = Describe("sendResponse", func() {
|
||||
@ -152,6 +154,24 @@ var _ = Describe("sendResponse", func() {
|
||||
})
|
||||
})
|
||||
|
||||
It("responds with HTTP 429 and Retry-After when the transcode limiter rejects", func() {
|
||||
w = httptest.NewRecorder()
|
||||
r = httptest.NewRequest("GET", "/rest/stream", nil)
|
||||
|
||||
sendError(w, r, fmt.Errorf("rejected: %w", stream.ErrTooManyTranscodes))
|
||||
|
||||
Expect(w.Code).To(Equal(http.StatusTooManyRequests))
|
||||
Expect(w.Header().Get("Retry-After")).ToNot(BeEmpty())
|
||||
|
||||
var subsonicResponse responses.Subsonic
|
||||
err := xml.Unmarshal(w.Body.Bytes(), &subsonicResponse)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(subsonicResponse.Status).To(Equal(responses.StatusFailed))
|
||||
Expect(subsonicResponse.Error).ToNot(BeNil())
|
||||
Expect(subsonicResponse.Error.Code).To(Equal(responses.ErrorGeneric))
|
||||
Expect(subsonicResponse.Error.Message).To(ContainSubstring("transcode"))
|
||||
})
|
||||
|
||||
It("updates status pointer when an error occurs", func() {
|
||||
pointer := int32(0)
|
||||
|
||||
|
||||
@ -1,12 +1,15 @@
|
||||
package subsonic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
"github.com/navidrome/navidrome/core/stream"
|
||||
"github.com/navidrome/navidrome/log"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/model/request"
|
||||
@ -119,14 +122,28 @@ func (api *Router) Download(w http.ResponseWriter, r *http.Request) (*responses.
|
||||
return nil, err
|
||||
case *model.Album:
|
||||
setHeaders(v.Name)
|
||||
return nil, api.archiver.ZipAlbum(ctx, id, format, maxBitRate, w)
|
||||
return nil, handleArchiveErr(ctx, id, api.archiver.ZipAlbum(ctx, id, format, maxBitRate, w))
|
||||
case *model.Artist:
|
||||
setHeaders(v.Name)
|
||||
return nil, api.archiver.ZipArtist(ctx, id, format, maxBitRate, w)
|
||||
return nil, handleArchiveErr(ctx, id, api.archiver.ZipArtist(ctx, id, format, maxBitRate, w))
|
||||
case *model.Playlist:
|
||||
setHeaders(v.Name)
|
||||
return nil, api.archiver.ZipPlaylist(ctx, id, format, maxBitRate, w)
|
||||
return nil, handleArchiveErr(ctx, id, api.archiver.ZipPlaylist(ctx, id, format, maxBitRate, w))
|
||||
default:
|
||||
return nil, model.ErrNotFound
|
||||
}
|
||||
}
|
||||
|
||||
// handleArchiveErr swallows ErrTooManyTranscodes from archive downloads so the
|
||||
// outer error handler does not try to write a 429 onto a response whose status
|
||||
// and Content-Disposition have already been flushed. The archive ends up with
|
||||
// the tracks that were written before the rejection (the rejected track and
|
||||
// any following ones are omitted); the server-side log is the unambiguous
|
||||
// signal operators can act on.
|
||||
func handleArchiveErr(ctx context.Context, id string, err error) error {
|
||||
if errors.Is(err, stream.ErrTooManyTranscodes) {
|
||||
log.Warn(ctx, "Archive download finalized early: transcode cap reached", "id", id, err)
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
@ -4,6 +4,7 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"maps"
|
||||
"net/http"
|
||||
"sync"
|
||||
"time"
|
||||
@ -76,9 +77,7 @@ func (t *requestThrottle) handler(next http.Handler) http.Handler {
|
||||
next.ServeHTTP(buf, r)
|
||||
}()
|
||||
|
||||
for k, v := range buf.header {
|
||||
w.Header()[k] = v
|
||||
}
|
||||
maps.Copy(w.Header(), buf.header)
|
||||
if buf.code > 0 {
|
||||
w.WriteHeader(buf.code)
|
||||
}
|
||||
|
||||
@ -272,18 +272,6 @@ const Player = () => {
|
||||
}
|
||||
}, [])
|
||||
|
||||
const onAudioSeeked = useCallback(
|
||||
(info) => {
|
||||
if (!info.isRadio && currentTrackId) {
|
||||
const posMs = Math.floor(info.currentTime * 1000)
|
||||
lastPositionMsRef.current = posMs
|
||||
const state = audioInstance?.paused ? 'paused' : 'playing'
|
||||
subsonic.reportPlayback(currentTrackId, posMs, state)
|
||||
}
|
||||
},
|
||||
[currentTrackId, audioInstance],
|
||||
)
|
||||
|
||||
const onAudioVolumeChange = useCallback(
|
||||
// sqrt to compensate for the logarithmic volume
|
||||
(volume) => dispatch(setVolume(Math.sqrt(volume))),
|
||||
@ -436,6 +424,35 @@ const Player = () => {
|
||||
}
|
||||
}, [isMobilePlayer, audioInstance])
|
||||
|
||||
// Report every seek (including programmatic ones the library does not surface
|
||||
// via onAudioSeeked, e.g. restartCurrentOnPrev). Debounce coalesces drag
|
||||
// bursts into one report at the final position.
|
||||
useEffect(() => {
|
||||
if (!audioInstance) return
|
||||
let timer = null
|
||||
const flush = () => {
|
||||
timer = null
|
||||
if (
|
||||
!currentTrackIdRef.current ||
|
||||
playerStateRef.current?.current?.isRadio
|
||||
) {
|
||||
return
|
||||
}
|
||||
const posMs = Math.floor((audioInstance.currentTime || 0) * 1000)
|
||||
const state = audioInstance.paused ? 'paused' : 'playing'
|
||||
subsonic.reportPlayback(currentTrackIdRef.current, posMs, state)
|
||||
}
|
||||
const handleSeeked = () => {
|
||||
if (timer) clearTimeout(timer)
|
||||
timer = setTimeout(flush, 250)
|
||||
}
|
||||
audioInstance.addEventListener('seeked', handleSeeked)
|
||||
return () => {
|
||||
if (timer) clearTimeout(timer)
|
||||
audioInstance.removeEventListener('seeked', handleSeeked)
|
||||
}
|
||||
}, [audioInstance])
|
||||
|
||||
return (
|
||||
<ThemeProvider theme={createMuiTheme(theme)}>
|
||||
<ReactJkMusicPlayer
|
||||
@ -444,7 +461,6 @@ const Player = () => {
|
||||
onAudioListsChange={onAudioListsChange}
|
||||
onAudioVolumeChange={onAudioVolumeChange}
|
||||
onAudioProgress={onAudioProgress}
|
||||
onAudioSeeked={onAudioSeeked}
|
||||
onAudioPlay={onAudioPlay}
|
||||
onAudioPlayTrackChange={onAudioPlayTrackChange}
|
||||
onAudioPause={onAudioPause}
|
||||
|
||||
@ -13,21 +13,10 @@ import { baseUrl, openInNewTab } from '../utils'
|
||||
import { httpClient } from '../dataProvider'
|
||||
|
||||
const Progress = (props) => {
|
||||
const { setLinked, setCheckingLink, apiKey } = props
|
||||
const { setLinked, setCheckingLink, openedTab } = props
|
||||
const notify = useNotify()
|
||||
let linkCheckDelay = 2000
|
||||
let linkChecks = 30
|
||||
const openedTab = useRef()
|
||||
|
||||
useEffect(() => {
|
||||
const callbackEndpoint = baseUrl(
|
||||
`/api/lastfm/link/callback?uid=${localStorage.getItem('userId')}`,
|
||||
)
|
||||
const callbackUrl = `${window.location.origin}${callbackEndpoint}`
|
||||
openedTab.current = openInNewTab(
|
||||
`https://www.last.fm/api/auth/?api_key=${apiKey}&cb=${callbackUrl}`,
|
||||
)
|
||||
}, [apiKey])
|
||||
|
||||
const endChecking = (success) => {
|
||||
linkCheckDelay = null
|
||||
@ -76,6 +65,7 @@ export const LastfmScrobbleToggle = (props) => {
|
||||
const [linked, setLinked] = useState(null)
|
||||
const [checkingLink, setCheckingLink] = useState(false)
|
||||
const [apiKey, setApiKey] = useState(false)
|
||||
const openedTab = useRef()
|
||||
|
||||
useEffect(() => {
|
||||
httpClient('/api/lastfm/link')
|
||||
@ -88,9 +78,42 @@ export const LastfmScrobbleToggle = (props) => {
|
||||
})
|
||||
}, [setLinked, setApiKey])
|
||||
|
||||
const startLink = () => {
|
||||
// Open the tab synchronously so popup blockers attribute it to the click.
|
||||
let tab
|
||||
try {
|
||||
tab = openInNewTab('about:blank')
|
||||
} catch {
|
||||
notify('message.lastfmLinkFailure', 'warning')
|
||||
return
|
||||
}
|
||||
openedTab.current = tab
|
||||
setCheckingLink(true)
|
||||
httpClient('/api/lastfm/link')
|
||||
.then((response) => {
|
||||
const linkToken = response.json.linkToken
|
||||
if (!linkToken) {
|
||||
tab?.close()
|
||||
notify('message.lastfmLinkFailure', 'warning')
|
||||
setCheckingLink(false)
|
||||
return
|
||||
}
|
||||
const callbackEndpoint = baseUrl(
|
||||
`/api/lastfm/link/callback?uid=${encodeURIComponent(linkToken)}`,
|
||||
)
|
||||
const callbackUrl = `${window.location.origin}${callbackEndpoint}`
|
||||
tab.location.href = `https://www.last.fm/api/auth/?api_key=${apiKey}&cb=${callbackUrl}`
|
||||
})
|
||||
.catch(() => {
|
||||
tab?.close()
|
||||
notify('message.lastfmLinkFailure', 'warning')
|
||||
setCheckingLink(false)
|
||||
})
|
||||
}
|
||||
|
||||
const toggleScrobble = () => {
|
||||
if (!linked) {
|
||||
setCheckingLink(true)
|
||||
startLink()
|
||||
} else {
|
||||
httpClient('/api/lastfm/link', { method: 'DELETE' })
|
||||
.then(() => {
|
||||
@ -121,7 +144,7 @@ export const LastfmScrobbleToggle = (props) => {
|
||||
<Progress
|
||||
setLinked={setLinked}
|
||||
setCheckingLink={setCheckingLink}
|
||||
apiKey={apiKey}
|
||||
openedTab={openedTab}
|
||||
/>
|
||||
)}
|
||||
{!apiKey && (
|
||||
|
||||
4
utils/cache/benchmark_test.go
vendored
4
utils/cache/benchmark_test.go
vendored
@ -116,7 +116,7 @@ func BenchmarkConcurrentCacheRead(b *testing.B) {
|
||||
for i := 0; i < b.N; i++ {
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(n)
|
||||
for g := 0; g < n; g++ {
|
||||
for range n {
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
s, err := fc.Get(context.Background(), item)
|
||||
@ -152,7 +152,7 @@ func BenchmarkConcurrentCacheMiss(b *testing.B) {
|
||||
wg.Add(n)
|
||||
// All goroutines request the SAME key (not yet cached)
|
||||
item := &benchItem{key: fmt.Sprintf("miss-%d", i)}
|
||||
for g := 0; g < n; g++ {
|
||||
for range n {
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
s, err := fc.Get(context.Background(), item)
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user