Merge 213afe325ee3ca807662dbf6f6212bbe2093851c into 600ea5482c36d3705fbca1d9ab749d9bdafd6f80

This commit is contained in:
TowyTowy 2026-07-30 22:47:45 +02:00 committed by GitHub
commit 656f04c621
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
2 changed files with 88 additions and 0 deletions

View File

@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"hash/fnv"
"io"
"mime"
"net/http"
@ -60,6 +61,20 @@ func (j *streamJob) Key() string {
return fmt.Sprintf("%s.%s.%d.%d.%d.%d.%s.%d", j.mf.ID, j.mf.UpdatedAt.Format(time.RFC3339Nano), j.bitRate, j.sampleRate, j.bitDepth, j.channels, j.format, j.offset)
}
// weakETag derives a weak HTTP entity tag from a transcode identity (the cache
// key). It is weak (W/) because ffmpeg output is not guaranteed to be
// byte-for-byte identical across runs (e.g. embedded timestamps/metadata), but
// it is a stable validator for a given representation: the same media file at
// the same bitrate/format/offset always yields the same tag, and it changes
// whenever any of those inputs change. Serving it lets a client revalidate a
// repeated request for the same transcoded track with a cheap 304 Not Modified
// instead of re-downloading the full audio.
func weakETag(identity string) string {
h := fnv.New64a()
_, _ = io.WriteString(h, identity)
return fmt.Sprintf(`W/"%x"`, h.Sum64())
}
// NewStream creates a Stream for the given MediaFile and Request. It handles both raw streaming (no transcoding)
// and transcoded streaming based on the requested format and bitrate. It also logs detailed information about
// the streaming request and whether the transcoding result was served from cache or not.
@ -109,6 +124,9 @@ func (ms *mediaStreamer) NewStream(ctx context.Context, mf *model.MediaFile, req
channels: req.Channels,
offset: req.Offset,
}
// Tag the transcoded representation so repeated requests for the same
// track within a playback can be revalidated (304) instead of re-fetched.
s.etag = weakETag(job.Key())
r, err := ms.cache.Get(ctx, job)
if err != nil {
// Rate-limit rejections are already logged at warn level by the
@ -137,6 +155,7 @@ type Stream struct {
mf *model.MediaFile
bitRate int
format string
etag string
io.ReadCloser
io.Seeker
}
@ -156,6 +175,14 @@ func (s *Stream) EstimatedContentLength() int {
// (meaning the HTTP 200 status has not been flushed yet and the caller can still send an error response).
// Empty output (0 bytes, no error) is logged but not treated as an error.
func (s *Stream) Serve(ctx context.Context, w http.ResponseWriter, r *http.Request) (int64, error) {
// Advertise a validator for the transcoded representation so clients can
// revalidate a repeated request with a conditional GET. On the seekable
// (fully cached) path http.ServeContent handles If-None-Match/If-Range
// against it, answering repeats with 304 instead of re-sending the audio.
if s.etag != "" {
w.Header().Set("ETag", s.etag)
}
if s.Seekable() {
http.ServeContent(w, r, s.Name(), s.ModTime(), s)
return -1, nil

View File

@ -4,6 +4,8 @@ import (
"context"
"errors"
"io"
"net/http"
"net/http/httptest"
"os"
"github.com/navidrome/navidrome/conf"
@ -138,5 +140,64 @@ var _ = Describe("MediaStreamer", func() {
Expect(err).To(BeNil())
Expect(s.Seekable()).To(BeTrue())
})
It("revalidates a repeated transcoded request with 304 instead of re-sending the audio", func() {
// Dedicated ffmpeg/cache so the transcoded payload is deterministic
// (the suite-level mock's reader is single-use and shared).
freshFFmpeg := tests.NewMockFFmpeg("some transcoded audio payload")
localCache := stream.NewTranscodingCache()
Eventually(func() bool { return localCache.Available(context.TODO()) }).Should(BeTrue())
localStreamer := stream.NewMediaStreamer(ds, freshFFmpeg, localCache)
// Prime the cache so the transcode is complete and seekable.
s, err := localStreamer.NewStream(ctx, mf, stream.Request{Format: "mp3", BitRate: 32})
Expect(err).ToNot(HaveOccurred())
_, _ = io.ReadAll(s)
_ = s.Close()
Eventually(func() bool { return freshFFmpeg.IsClosed() }, "3s").Should(BeTrue())
// First request: serve fully and capture the advertised ETag.
s, err = localStreamer.NewStream(ctx, mf, stream.Request{Format: "mp3", BitRate: 32})
Expect(err).ToNot(HaveOccurred())
Expect(s.Seekable()).To(BeTrue())
rec := httptest.NewRecorder()
_, err = s.Serve(ctx, rec, httptest.NewRequest(http.MethodGet, "/stream", nil))
Expect(err).ToNot(HaveOccurred())
_ = s.Close()
Expect(rec.Code).To(Equal(http.StatusOK))
Expect(rec.Body.Len()).To(BeNumerically(">", 0))
etag := rec.Header().Get("ETag")
Expect(etag).To(HavePrefix(`W/"`))
// Repeat with a matching If-None-Match: must be a bodiless 304.
s, err = localStreamer.NewStream(ctx, mf, stream.Request{Format: "mp3", BitRate: 32})
Expect(err).ToNot(HaveOccurred())
rec2 := httptest.NewRecorder()
r2 := httptest.NewRequest(http.MethodGet, "/stream", nil)
r2.Header.Set("If-None-Match", etag)
_, err = s.Serve(ctx, rec2, r2)
Expect(err).ToNot(HaveOccurred())
_ = s.Close()
Expect(rec2.Code).To(Equal(http.StatusNotModified))
Expect(rec2.Body.Len()).To(Equal(0))
})
It("changes the ETag when the requested bitrate changes", func() {
s1, err := streamer.NewStream(ctx, mf, stream.Request{Format: "mp3", BitRate: 32})
Expect(err).ToNot(HaveOccurred())
rec1 := httptest.NewRecorder()
_, _ = s1.Serve(ctx, rec1, httptest.NewRequest(http.MethodGet, "/stream", nil))
_ = s1.Close()
s2, err := streamer.NewStream(ctx, mf, stream.Request{Format: "mp3", BitRate: 64})
Expect(err).ToNot(HaveOccurred())
rec2 := httptest.NewRecorder()
_, _ = s2.Serve(ctx, rec2, httptest.NewRequest(http.MethodGet, "/stream", nil))
_ = s2.Close()
Expect(rec1.Header().Get("ETag")).ToNot(BeEmpty())
Expect(rec2.Header().Get("ETag")).ToNot(BeEmpty())
Expect(rec1.Header().Get("ETag")).ToNot(Equal(rec2.Header().Get("ETag")))
})
})
})