mirror of
https://github.com/navidrome/navidrome.git
synced 2026-08-31 07:30:32 +00:00
* fix(stream): abort the response when a transcoded stream is truncated When a transcode failed after some audio had already been sent, Serve logged the error and returned nil, so Go finished the chunked body normally and the client received an apparently complete, silently short file. Symfonium users hit this on large offline syncs, and the worst path, ffmpeg dying mid-write behind the transcoding cache, produced no error and nothing in the log above Debug: the cache writer was closed plainly, so readers drained the truncated entry to a clean EOF. The root cause of that silence is an fscache limitation: Close is the only way to end a cache write, and Close always means "complete". This adopts the deluan/fscache fork, which adds CloseWithError: on failure copyAndClose now cancels the entry with the cause, so every attached reader fails mid-read with the real error instead of EOF, a late Get for the entry is refused, and the entry never reports a final size. The error travels inside the entry each reader holds, which makes per-generation delivery automatic and needs no bookkeeping on our side. With the failure arriving in-band, one change in Serve covers every mode: an io.Copy error after bytes are on the wire panics with http.ErrAbortHandler. Go aborts the response without the terminating chunk (RST_STREAM on HTTP/2), chi's Recoverer re-panics that value, and the deferred stream.Close() still runs, so the transcode limiter slot is released as before. Two behaviors improve as side effects. A transcoder that dies before its first byte now yields a Subsonic error response instead of a 200 with an empty body, since the failure reaches Serve as an error while the status is still unsent; genuinely empty output (clean EOF, exit 0) keeps the 200. And a failed entry's invalidation no longer defers its unlink past a replacement entry re-creating the same file, because canceling already closed its readers. * fix(cache): warn when the cache writer cannot report failures to readers The CloseWithError capability comes from the fscache fork via a go.mod replace directive, and a type assertion picks it up. If that directive is ever lost, the assertion fails silently, readers of a dead writer go back to draining a truncated entry to a clean EOF, and nothing says so. Two layers against that: a warning on the failure path when the writer lacks the capability, and a test that asserts the writer fscache returns carries it, so losing the fork fails CI instead of a listener's download. * build: point the fscache replace at the fork's master deluan/fscache#1 is merged; pin the merge commit instead of the review branch. Pinned by sha because the module proxy still resolves the fork's master ref to its pre-merge commit. * build: reference the upstream fscache PR in the replace comment The replace itself must keep pointing at the fork: the commit only exists in djherbis/fscache under refs/pull/22/head, which the Go module fetcher cannot resolve (verified: unknown revision for both short and full sha). The same commit is advertised on the fork's master, so that is the fetchable source.
446 lines
16 KiB
Go
446 lines
16 KiB
Go
package cache
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/navidrome/navidrome/conf"
|
|
"github.com/navidrome/navidrome/conf/configtest"
|
|
. "github.com/onsi/ginkgo/v2"
|
|
. "github.com/onsi/gomega"
|
|
)
|
|
|
|
// Call NewFileCache and wait for it to be ready
|
|
func callNewFileCache(name, cacheSize, cacheFolder string, maxItems int, getReader ReadFunc) *fileCache {
|
|
fc := NewFileCache(name, cacheSize, cacheFolder, maxItems, getReader).(*fileCache)
|
|
Eventually(func() bool { return fc.ready.Load() }, 10*time.Second).Should(BeTrue())
|
|
return fc
|
|
}
|
|
|
|
var _ = Describe("File Caches", func() {
|
|
BeforeEach(func() {
|
|
tmpDir, _ := os.MkdirTemp("", "file_caches")
|
|
DeferCleanup(func() {
|
|
configtest.SetupConfig()
|
|
_ = os.RemoveAll(tmpDir)
|
|
})
|
|
conf.Server.CacheFolder = conf.NewDir(tmpDir)
|
|
})
|
|
|
|
Describe("NewFileCache", func() {
|
|
It("creates the cache folder", func() {
|
|
Expect(callNewFileCache("test", "1k", "test", 0, nil)).ToNot(BeNil())
|
|
|
|
_, err := os.Stat(filepath.Join(conf.Server.CacheFolder.String(), "test"))
|
|
Expect(os.IsNotExist(err)).To(BeFalse())
|
|
})
|
|
|
|
It("creates the cache folder with invalid size", func() {
|
|
fc := callNewFileCache("test", "abc", "test", 0, nil)
|
|
Expect(fc.cache).ToNot(BeNil())
|
|
Expect(fc.disabled).To(BeFalse())
|
|
})
|
|
|
|
It("returns empty if cache size is '0'", func() {
|
|
fc := callNewFileCache("test", "0", "test", 0, nil)
|
|
Expect(fc.cache).To(BeNil())
|
|
Expect(fc.disabled).To(BeTrue())
|
|
})
|
|
|
|
It("reports when cache is disabled", func() {
|
|
fc := callNewFileCache("test", "0", "test", 0, nil)
|
|
Expect(fc.Disabled(context.Background())).To(BeTrue())
|
|
fc = callNewFileCache("test", "1KB", "test", 0, nil)
|
|
Expect(fc.Disabled(context.Background())).To(BeFalse())
|
|
})
|
|
})
|
|
|
|
Describe("FileCache", func() {
|
|
It("caches data if cache is enabled", func() {
|
|
called := false
|
|
fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
called = true
|
|
return strings.NewReader(arg.Key()), nil
|
|
})
|
|
// First call is a MISS
|
|
s, err := fc.Get(context.Background(), &testArg{"test"})
|
|
Expect(err).To(BeNil())
|
|
Expect(s.Cached).To(BeFalse())
|
|
Expect(s.Closer).To(BeNil())
|
|
Expect(io.ReadAll(s)).To(Equal([]byte("test")))
|
|
|
|
// Second call is a HIT
|
|
called = false
|
|
s, err = fc.Get(context.Background(), &testArg{"test"})
|
|
Expect(err).To(BeNil())
|
|
Expect(io.ReadAll(s)).To(Equal([]byte("test")))
|
|
Expect(s.Cached).To(BeTrue())
|
|
Expect(s.Closer).ToNot(BeNil())
|
|
Expect(called).To(BeFalse())
|
|
})
|
|
|
|
It("does not cache data if cache is disabled", func() {
|
|
called := false
|
|
fc := callNewFileCache("test", "0", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
called = true
|
|
return strings.NewReader(arg.Key()), nil
|
|
})
|
|
// First call is a MISS
|
|
s, err := fc.Get(context.Background(), &testArg{"test"})
|
|
Expect(err).To(BeNil())
|
|
Expect(s.Cached).To(BeFalse())
|
|
Expect(io.ReadAll(s)).To(Equal([]byte("test")))
|
|
|
|
// Second call is also a MISS
|
|
called = false
|
|
s, err = fc.Get(context.Background(), &testArg{"test"})
|
|
Expect(err).To(BeNil())
|
|
Expect(io.ReadAll(s)).To(Equal([]byte("test")))
|
|
Expect(s.Cached).To(BeFalse())
|
|
Expect(called).To(BeTrue())
|
|
})
|
|
|
|
It("writes a completion marker after a successful cache write", func() {
|
|
fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return strings.NewReader("complete-data"), nil
|
|
})
|
|
s, err := fc.Get(context.Background(), &testArg{"markme"})
|
|
Expect(err).To(BeNil())
|
|
_, _ = io.ReadAll(s)
|
|
_ = s.Close()
|
|
|
|
// EOF must imply the entry is settled on disk (Windows temp-dir cleanups rely on it).
|
|
dataPath := fcSpreadFS(fc).KeyMapper((&testArg{"markme"}).Key())
|
|
_, statErr := os.Stat(dataPath + ".complete")
|
|
Expect(statErr).ToNot(HaveOccurred())
|
|
})
|
|
|
|
It("serves a concurrent reader from an in-progress write and marks complete once", func() {
|
|
pr, pw := io.Pipe()
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return pr, nil // slow, still-being-produced stream
|
|
})
|
|
|
|
// First Get → MISS; the cache starts copying pr into the entry in a goroutine.
|
|
s1, err := fc.Get(context.Background(), &testArg{"live"})
|
|
Expect(err).To(BeNil())
|
|
|
|
// Write the first chunk so the entry exists with in-flight bytes,
|
|
// but leave the pipe open so the second reader can attach mid-stream.
|
|
// io.Pipe writes block until the cache goroutine reads them, giving us
|
|
// a deterministic happens-before: the entry is live before we call Get again.
|
|
_, err = pw.Write([]byte("hello "))
|
|
Expect(err).To(BeNil())
|
|
|
|
// Second Get while the pipe is still open → attaches to the in-progress entry.
|
|
s2, err := fc.Get(context.Background(), &testArg{"live"})
|
|
Expect(err).To(BeNil())
|
|
|
|
// Drain both readers concurrently; they race against the producer below.
|
|
ch1 := make(chan []byte, 1)
|
|
ch2 := make(chan []byte, 1)
|
|
go func() { b, _ := io.ReadAll(s1); ch1 <- b }()
|
|
go func() { b, _ := io.ReadAll(s2); ch2 <- b }()
|
|
|
|
// Deliver the rest of the stream and close; both draining goroutines must see it.
|
|
_, err = pw.Write([]byte("world"))
|
|
Expect(err).To(BeNil())
|
|
Expect(pw.Close()).To(Succeed())
|
|
|
|
Expect(string(<-ch1)).To(Equal("hello world"))
|
|
Expect(string(<-ch2)).To(Equal("hello world"))
|
|
_ = s1.Close()
|
|
_ = s2.Close()
|
|
|
|
// Exactly one completion marker must appear.
|
|
dataPath := fcSpreadFS(fc).KeyMapper((&testArg{"live"}).Key())
|
|
Eventually(func() bool {
|
|
_, e := os.Stat(dataPath + ".complete")
|
|
return e == nil
|
|
}).Should(BeTrue())
|
|
|
|
// Steady-state HIT: full data, Cached flag set.
|
|
s3, err := fc.Get(context.Background(), &testArg{"live"})
|
|
Expect(err).To(BeNil())
|
|
got3, _ := io.ReadAll(s3)
|
|
_ = s3.Close()
|
|
Expect(s3.Cached).To(BeTrue())
|
|
Expect(string(got3)).To(Equal("hello world"))
|
|
})
|
|
|
|
Context("reader errors", func() {
|
|
When("creating a reader fails", func() {
|
|
It("does not cache", func() {
|
|
fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return nil, errors.New("failed")
|
|
})
|
|
|
|
_, err := fc.Get(context.Background(), &testArg{"test"})
|
|
Expect(err).To(MatchError("failed"))
|
|
})
|
|
})
|
|
When("reader returns error", func() {
|
|
It("does not cache", func() {
|
|
fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return errFakeReader{errors.New("read failure")}, nil
|
|
})
|
|
|
|
s, err := fc.Get(context.Background(), &testArg{"test"})
|
|
Expect(err).ToNot(HaveOccurred())
|
|
_, _ = io.Copy(io.Discard, s)
|
|
// TODO How to make the fscache reader return the underlying reader error?
|
|
//Expect(err).To(MatchError("read failure"))
|
|
|
|
// Data should not be cached (or eventually be removed from cache)
|
|
Eventually(func() bool {
|
|
s, _ = fc.Get(context.Background(), &testArg{"test"})
|
|
if s != nil {
|
|
return s.Cached
|
|
}
|
|
return false
|
|
}).Should(BeFalse())
|
|
})
|
|
})
|
|
})
|
|
|
|
Context("crash leftover (issue #5636)", func() {
|
|
It("does not serve a partial file left on disk as a complete HIT", func() {
|
|
// First init: empties + writes the migration sentinel.
|
|
fc1 := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return strings.NewReader("UNUSED"), nil
|
|
})
|
|
_ = fc1
|
|
|
|
// Plant a partial file (no marker), simulating a killed process.
|
|
sfs, err := NewSpreadFS(filepath.Join(conf.Server.CacheFolder.String(), "test"), 0755)
|
|
Expect(err).To(BeNil())
|
|
partialPath := sfs.KeyMapper((&testArg{"track"}).Key())
|
|
Expect(os.MkdirAll(filepath.Dir(partialPath), 0755)).To(Succeed())
|
|
Expect(os.WriteFile(partialPath, []byte("PARTIAL"), 0600)).To(Succeed())
|
|
|
|
// "Restart": a fresh cache over the same folder (sentinel present → strict).
|
|
getReaderCalled := false
|
|
fc2 := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
getReaderCalled = true
|
|
return strings.NewReader("FULL-TRANSCODE"), nil
|
|
})
|
|
|
|
s, err := fc2.Get(context.Background(), &testArg{"track"})
|
|
Expect(err).To(BeNil())
|
|
data, _ := io.ReadAll(s)
|
|
_ = s.Close()
|
|
|
|
Expect(getReaderCalled).To(BeTrue()) // re-transcoded, not served stale
|
|
Expect(string(data)).To(Equal("FULL-TRANSCODE"))
|
|
})
|
|
})
|
|
|
|
Context("live error path still invalidates", func() {
|
|
It("leaves no data file and no marker after a mid-stream reader error", func() {
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return errFakeReader{errors.New("boom")}, nil
|
|
})
|
|
s, err := fc.Get(context.Background(), &testArg{"err"})
|
|
Expect(err).To(BeNil())
|
|
_, _ = io.Copy(io.Discard, s)
|
|
_ = s.Close()
|
|
|
|
dataPath := fcSpreadFS(fc).KeyMapper((&testArg{"err"}).Key())
|
|
Eventually(func() bool {
|
|
_, e1 := os.Stat(dataPath)
|
|
_, e2 := os.Stat(dataPath + ".complete")
|
|
return os.IsNotExist(e1) && os.IsNotExist(e2)
|
|
}).Should(BeTrue())
|
|
})
|
|
|
|
It("gets a writer that can report failures to readers", func() {
|
|
// Guards the fork adoption: if the fscache replace directive is ever lost,
|
|
// this fails in CI instead of silently reviving the truncation bug.
|
|
fc := callNewFileCache("test", "10MB", "test", 0, nil)
|
|
_, w, err := fc.cache.Get("capability")
|
|
Expect(err).To(BeNil())
|
|
DeferCleanup(func() { _ = w.Close() })
|
|
|
|
_, ok := w.(interface{ CloseWithError(error) error })
|
|
Expect(ok).To(BeTrue(), "fscache writer lost CloseWithError; check the go.mod replace directive")
|
|
})
|
|
|
|
It("fails the reader with the cause instead of a clean EOF", func() {
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return &partialThenErrReader{data: []byte("PARTIAL"), err: errors.New("transcoder died")}, nil
|
|
})
|
|
s, err := fc.Get(context.Background(), &testArg{"inband"})
|
|
Expect(err).To(BeNil())
|
|
DeferCleanup(func() { _ = s.Close() })
|
|
|
|
_, err = io.ReadAll(s)
|
|
Expect(err).To(MatchError(ContainSubstring("transcoder died")))
|
|
})
|
|
|
|
It("fails a reader that joined mid-write with the same cause", func() {
|
|
pr, pw := io.Pipe()
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return pr, nil
|
|
})
|
|
s1, err := fc.Get(context.Background(), &testArg{"joined"})
|
|
Expect(err).To(BeNil())
|
|
DeferCleanup(func() { _ = s1.Close() })
|
|
|
|
// The blocking pipe write gives a happens-before: the entry is in flight.
|
|
_, err = pw.Write([]byte("PARTIAL"))
|
|
Expect(err).To(BeNil())
|
|
|
|
s2, err := fc.Get(context.Background(), &testArg{"joined"})
|
|
Expect(err).To(BeNil())
|
|
DeferCleanup(func() { _ = s2.Close() })
|
|
Expect(s2.Cached).To(BeTrue())
|
|
|
|
Expect(pw.CloseWithError(errors.New("transcoder died"))).To(Succeed())
|
|
|
|
_, err = io.ReadAll(s2)
|
|
Expect(err).To(MatchError(ContainSubstring("transcoder died")))
|
|
})
|
|
|
|
It("does not write a completion marker when the write fails after partial bytes", func() {
|
|
// Mimics a transcode that produces real output and then dies:
|
|
// the bytes land on disk, but the entry must NOT be marked complete.
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
return &partialThenErrReader{data: []byte("PARTIAL-OUTPUT"), err: errors.New("transcoder died")}, nil
|
|
})
|
|
s, err := fc.Get(context.Background(), &testArg{"partial"})
|
|
Expect(err).To(BeNil())
|
|
_, _ = io.Copy(io.Discard, s)
|
|
_ = s.Close()
|
|
|
|
dataPath := fcSpreadFS(fc).KeyMapper((&testArg{"partial"}).Key())
|
|
// The marker must never appear for a failed write. Give the async
|
|
// writer time to finish, then assert the marker stays absent.
|
|
Consistently(func() bool {
|
|
_, e := os.Stat(dataPath + ".complete")
|
|
return os.IsNotExist(e)
|
|
}).Should(BeTrue())
|
|
})
|
|
})
|
|
|
|
Context("entry outliving its data file", func() {
|
|
It("re-fetches when the data file vanished behind the cache's back", func() {
|
|
var calls atomic.Int32
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
calls.Add(1)
|
|
return strings.NewReader("payload"), nil
|
|
})
|
|
|
|
s, err := fc.Get(context.Background(), &testArg{"vanish"})
|
|
Expect(err).To(BeNil())
|
|
Expect(io.ReadAll(s)).To(Equal([]byte("payload")))
|
|
Expect(s.Close()).To(Succeed())
|
|
|
|
dataPath := fcSpreadFS(fc).KeyMapper((&testArg{"vanish"}).Key())
|
|
Eventually(func() error { _, e := os.Stat(dataPath); return e }).Should(Succeed())
|
|
Expect(os.Remove(dataPath)).To(Succeed())
|
|
|
|
s2, err := fc.Get(context.Background(), &testArg{"vanish"})
|
|
Expect(err).ToNot(HaveOccurred())
|
|
Expect(io.ReadAll(s2)).To(Equal([]byte("payload")))
|
|
_ = s2.Close()
|
|
Expect(calls.Load()).To(BeNumerically("==", 2))
|
|
})
|
|
|
|
It("removes a failed entry promptly, without eating its replacement", func() {
|
|
// Cancel closes the failed entry's readers, so its removal no longer defers
|
|
// past the point where a new entry re-creates the same file.
|
|
var n atomic.Int32
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
if n.Add(1) == 1 {
|
|
return &partialThenErrReader{data: []byte("PARTIAL"), err: errors.New("died")}, nil
|
|
}
|
|
return strings.NewReader("GOOD"), nil
|
|
})
|
|
|
|
key := (&testArg{"deferred"}).Key()
|
|
s1, err := fc.Get(context.Background(), &testArg{"deferred"})
|
|
Expect(err).To(BeNil())
|
|
|
|
Eventually(func() bool { return fc.cache.Exists(key) }).Should(BeFalse())
|
|
|
|
s2, err := fc.Get(context.Background(), &testArg{"deferred"})
|
|
Expect(err).To(BeNil())
|
|
Expect(io.ReadAll(s2)).To(Equal([]byte("GOOD")))
|
|
Expect(s2.Close()).To(Succeed())
|
|
|
|
Expect(s1.Close()).To(Succeed())
|
|
|
|
dataPath := fcSpreadFS(fc).KeyMapper(key)
|
|
Consistently(func() error {
|
|
_, e := os.Stat(dataPath)
|
|
return e
|
|
}).Should(Succeed(), "the replacement entry's file must survive the failed entry's cleanup")
|
|
|
|
s3, err := fc.Get(context.Background(), &testArg{"deferred"})
|
|
Expect(err).ToNot(HaveOccurred())
|
|
Expect(io.ReadAll(s3)).To(Equal([]byte("GOOD")))
|
|
_ = s3.Close()
|
|
Expect(n.Load()).To(Equal(int32(2)), "the third Get must be served from cache")
|
|
})
|
|
|
|
It("re-fetches when an adopted entry's data file vanished", func() {
|
|
// Entries adopted on startup take a different code path than in-process ones.
|
|
var calls atomic.Int32
|
|
fc := callNewFileCache("test", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
calls.Add(1)
|
|
return strings.NewReader("payload"), nil
|
|
})
|
|
dataPath := fcSpreadFS(fc).KeyMapper((&testArg{"adopted"}).Key())
|
|
Expect(os.MkdirAll(filepath.Dir(dataPath), 0755)).To(Succeed())
|
|
Expect(os.WriteFile(dataPath, []byte("payload"), 0600)).To(Succeed())
|
|
Expect(fcSpreadFS(fc).MarkComplete(dataPath)).To(Succeed())
|
|
|
|
adopted := callNewFileCache("test2", "10MB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) {
|
|
calls.Add(1)
|
|
return strings.NewReader("payload"), nil
|
|
})
|
|
Expect(os.Remove(dataPath)).To(Succeed())
|
|
|
|
s, err := adopted.Get(context.Background(), &testArg{"adopted"})
|
|
Expect(err).ToNot(HaveOccurred())
|
|
Expect(io.ReadAll(s)).To(Equal([]byte("payload")))
|
|
_ = s.Close()
|
|
})
|
|
})
|
|
})
|
|
})
|
|
|
|
type testArg struct{ s string }
|
|
|
|
func (t *testArg) Key() string { return t.s }
|
|
|
|
type errFakeReader struct{ err error }
|
|
|
|
func (e errFakeReader) Read([]byte) (int, error) { return 0, e.err }
|
|
|
|
// partialThenErrReader emits data once, then fails — mimicking a transcoder
|
|
// that produces some output and then dies mid-stream.
|
|
type partialThenErrReader struct {
|
|
data []byte
|
|
err error
|
|
done bool
|
|
}
|
|
|
|
func (r *partialThenErrReader) Read(p []byte) (int, error) {
|
|
if r.done {
|
|
return 0, r.err
|
|
}
|
|
r.done = true
|
|
return copy(p, r.data), nil
|
|
}
|
|
|
|
func fcSpreadFS(fc *fileCache) *spreadFS {
|
|
return fc.fs
|
|
}
|