From 95f67d2c4ef391967327f0c7dccd5d0d8b02f62e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Deluan=20Quint=C3=A3o?= Date: Sat, 3 Oct 2026 08:58:54 -0700 Subject: [PATCH] fix(share): reuse cached transcodes for share streams and zip downloads (#6262) * fix(share): reuse cached transcodes when streaming from share links Public share streams built the stream request with only the share's format and bit rate, leaving sample rate, bit depth and channels at zero. Regular playback resolves those through the transcode decider (e.g. 48000 Hz for Opus), and they are part of the transcoding cache key, so a track already transcoded during normal playback was transcoded again into a separate, identical cache entry when played through a share link. The public router now resolves share stream requests with the same TranscodeDecider.ResolveRequest used by the Subsonic stream endpoint, so both paths produce the same request and share cache entries. Fixes #6261 * fix(archiver): reuse cached transcodes when zipping downloads Zip downloads (album, artist, playlist and share) built the stream request with only the format and bit rate, leaving sample rate, bit depth and channels at zero. Those are part of the transcoding cache key, so a track already transcoded for playback was transcoded again into a separate cache entry when downloaded in a zip, and vice versa. The archiver now resolves each request with TranscodeDecider.ResolveRequest, the same as single-song downloads and streams. This also applies the decider's defaults, so a zip requested without a bit rate uses the target format's default bit rate instead of leaving it to ffmpeg. * fix(archiver): name zip entries after the resolved transcoding format The transcode decider can pick a different format than the one requested (for example a player's forced transcoding, or a fallback to the default downsampling format when the requested one can't be produced). Zip entry names and the playlist M3U were still built from the requested format, so an entry could end in .mp3 or .flac while holding Opus data. Each track's request is now resolved before its entry name is built, and the name uses the resolved format. --- cmd/wire_gen.go | 9 +-- core/archiver.go | 30 ++++++---- core/archiver_test.go | 58 +++++++++++++++++++- server/public/handle_streams.go | 5 +- server/public/handle_streams_test.go | 17 ++++-- server/public/public.go | 5 +- server/subsonic/e2e/subsonic_artwork_test.go | 2 +- 7 files changed, 101 insertions(+), 25 deletions(-) diff --git a/cmd/wire_gen.go b/cmd/wire_gen.go index 19f92d9d5..2a396689d 100644 --- a/cmd/wire_gen.go +++ b/cmd/wire_gen.go @@ -95,8 +95,9 @@ func CreateSubsonicAPIRouter(ctx context.Context) *subsonic.Router { artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg) transcodingCache := stream.GetTranscodingCache() mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) + transcodeDecider := stream.NewTranscodeDecider(dataStore, fFmpeg) share := core.NewShare(dataStore) - archiver := core.NewArchiver(mediaStreamer, dataStore, share, artworkArtwork) + archiver := core.NewArchiver(mediaStreamer, transcodeDecider, dataStore, share, artworkArtwork) players := core.NewPlayers(dataStore) broker := events.GetBroker() metricsMetrics := metrics.GetPrometheusInstance(dataStore) @@ -110,7 +111,6 @@ func CreateSubsonicAPIRouter(ctx context.Context) *subsonic.Router { playTracker := scrobbler.GetPlayTracker(dataStore, broker, manager) playbackServer := playback.GetInstance(dataStore) lyricsLyrics := lyrics.NewLyrics(dataStore, manager) - transcodeDecider := stream.NewTranscodeDecider(dataStore, fFmpeg) sonicSonic := sonic.New(dataStore, manager, matcherMatcher) router := subsonic.New(dataStore, artworkArtwork, mediaStreamer, archiver, players, provider, modelScanner, broker, playlistsPlaylists, playTracker, share, playbackServer, metricsMetrics, lyricsLyrics, transcodeDecider, sonicSonic) return router @@ -159,9 +159,10 @@ func CreatePublicRouter() *public.Router { artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg) transcodingCache := stream.GetTranscodingCache() mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) + transcodeDecider := stream.NewTranscodeDecider(dataStore, fFmpeg) share := core.NewShare(dataStore) - archiver := core.NewArchiver(mediaStreamer, dataStore, share, artworkArtwork) - router := public.New(dataStore, artworkArtwork, mediaStreamer, share, archiver) + archiver := core.NewArchiver(mediaStreamer, transcodeDecider, dataStore, share, artworkArtwork) + router := public.New(dataStore, artworkArtwork, mediaStreamer, transcodeDecider, share, archiver) return router } diff --git a/core/archiver.go b/core/archiver.go index 6f362322a..60eb44858 100644 --- a/core/archiver.go +++ b/core/archiver.go @@ -35,13 +35,14 @@ type Archiver interface { ZipPlaylist(ctx context.Context, id string, format string, bitrate int, w io.Writer) error } -func NewArchiver(ms stream.MediaStreamer, ds model.DataStore, shares Share, artwork artwork.Artwork) Archiver { - return &archiver{ds: ds, ms: ms, shares: shares, artwork: artwork} +func NewArchiver(ms stream.MediaStreamer, decider stream.TranscodeDecider, ds model.DataStore, shares Share, artwork artwork.Artwork) Archiver { + return &archiver{ds: ds, ms: ms, decider: decider, shares: shares, artwork: artwork} } type archiver struct { ds model.DataStore ms stream.MediaStreamer + decider stream.TranscodeDecider shares Share artwork artwork.Artwork } @@ -78,8 +79,9 @@ func (a *archiver) zipAlbums(ctx context.Context, id string, format string, bitr log.Debug(ctx, "Zipping album", "name", album[0].Album, "artist", album[0].AlbumArtist, "folder", folder, "format", format, "bitrate", bitrate, "isMultiDisc", isMultiDisc, "numTracks", len(album)) for _, mf := range album { - file := a.albumFilename(mf, format, isMultiDisc, folder) - if addErr := a.addFileToZip(ctx, z, mf, format, bitrate, file); errors.Is(addErr, stream.ErrTooManyTranscodes) { + req := a.resolveRequest(ctx, &mf, format, bitrate) + file := a.albumFilename(mf, req.Format, isMultiDisc, folder) + if addErr := a.addFileToZip(ctx, z, mf, req, 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 @@ -204,8 +206,9 @@ 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) - if addErr := a.addFileToZip(ctx, z, mf, format, bitrate, file); errors.Is(addErr, stream.ErrTooManyTranscodes) { + req := a.resolveRequest(ctx, &mf, format, bitrate) + file := a.playlistFilename(mf, req.Format, idx) + if addErr := a.addFileToZip(ctx, z, mf, req, 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() @@ -251,7 +254,14 @@ func (a *archiver) playlistFilename(mf model.MediaFile, format string, idx int) return fmt.Sprintf("%02d - %s - %s.%s", idx+1, str.SanitizeFilename(mf.Artist), str.SanitizeFilename(mf.Title), ext) } -func (a *archiver) addFileToZip(ctx context.Context, z *zip.Writer, mf model.MediaFile, format string, bitrate int, filename string) error { +func (a *archiver) resolveRequest(ctx context.Context, mf *model.MediaFile, format string, bitrate int) stream.Request { + if format == "" || format == "raw" { + return stream.Request{Format: "raw"} + } + return a.decider.ResolveRequest(ctx, mf, format, bitrate, 0) +} + +func (a *archiver) addFileToZip(ctx context.Context, z *zip.Writer, mf model.MediaFile, req stream.Request, filename string) error { path := mf.AbsolutePath() // Open the source before writing the zip entry header so a rejection @@ -259,13 +269,13 @@ func (a *archiver) addFileToZip(ctx context.Context, z *zip.Writer, mf model.Med // 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}) + if req.Format != "raw" { + r, err = a.ms.NewStream(ctx, &mf, req) } else { r, err = os.Open(path) } if err != nil { - log.Error(ctx, "Error opening file for zipping", "file", path, "format", format, err) + log.Error(ctx, "Error opening file for zipping", "file", path, "format", req.Format, err) return err } defer func() { diff --git a/core/archiver_test.go b/core/archiver_test.go index 4e00ce78c..178d1b6b9 100644 --- a/core/archiver_test.go +++ b/core/archiver_test.go @@ -26,6 +26,7 @@ var _ = Describe("Archiver", func() { var ( arch core.Archiver ms *mockMediaStreamer + dc *fakeDecider ds *mockDataStore sh *mockShare ca *mockCoverArt @@ -33,10 +34,11 @@ var _ = Describe("Archiver", func() { BeforeEach(func() { ms = &mockMediaStreamer{} + dc = &fakeDecider{} sh = &mockShare{} ds = &mockDataStore{} ca = &mockCoverArt{images: map[string][]byte{}} - arch = core.NewArchiver(ms, ds, sh, ca) + arch = core.NewArchiver(ms, dc, ds, sh, ca) }) Context("ZipAlbum", func() { @@ -66,6 +68,23 @@ var _ = Describe("Archiver", func() { Expect(zr.File[0].Name).To(Equal("Album_Promo/01 - track1.mp3")) Expect(zr.File[1].Name).To(Equal("Album_Promo/02 - track2.mp3")) }) + + It("streams the request resolved by the transcode decider and names the entry after its format", func() { + mfRepo := &mockMediaFileRepository{} + mfRepo.On("GetAll", mock.Anything).Return(model.MediaFiles{{Path: "test_data/01 - track1.flac", Suffix: "flac", AlbumID: "1"}}, nil) + ds.On("MediaFile").Return(mfRepo) + resolved := stream.Request{Format: "opus", BitRate: 128, SampleRate: 48000, Channels: 2} + dc.resolved = &resolved + ms.On("NewStream", mock.Anything, mock.Anything, resolved).Return(io.NopCloser(strings.NewReader("test")), nil).Once() + + out := new(bytes.Buffer) + Expect(arch.ZipAlbum(GinkgoT().Context(), "1", "mp3", 128, out)).To(Succeed()) + ms.AssertExpectations(GinkgoT()) + + zr, err := zip.NewReader(bytes.NewReader(out.Bytes()), int64(out.Len())) + Expect(err).ToNot(HaveOccurred()) + Expect(zr.File[0].Name).To(HaveSuffix("01 - track1.opus")) + }) }) Context("ZipArtist", func() { @@ -296,6 +315,30 @@ var _ = Describe("Archiver", func() { }) Context("ZipPlaylist", func() { + It("names the entries and the M3U lines after the resolved format", func() { + pls := &model.Playlist{ID: "1", Name: "Test Playlist", Tracks: []model.PlaylistTrack{ + {MediaFile: model.MediaFile{Path: "test_data/01 - track1.flac", Suffix: "flac", Artist: "Artist 1", Title: "track1"}}, + }} + plRepo := &mockPlaylistRepository{} + plRepo.On("GetWithTracks", "1", true, false).Return(pls, nil) + ds.On("Playlist").Return(plRepo) + dc.resolved = &stream.Request{Format: "opus", BitRate: 128} + ms.On("NewStream", mock.Anything, mock.Anything, *dc.resolved).Return(io.NopCloser(strings.NewReader("test")), nil) + + out := new(bytes.Buffer) + Expect(arch.ZipPlaylist(GinkgoT().Context(), "1", "mp3", 128, out)).To(Succeed()) + + zr, err := zip.NewReader(bytes.NewReader(out.Bytes()), int64(out.Len())) + Expect(err).ToNot(HaveOccurred()) + Expect(zr.File[0].Name).To(Equal("01 - Artist 1 - track1.opus")) + m3u, err := zr.File[1].Open() + Expect(err).ToNot(HaveOccurred()) + defer m3u.Close() + content, err := io.ReadAll(m3u) + Expect(err).ToNot(HaveOccurred()) + Expect(string(content)).To(ContainSubstring("01 - Artist 1 - track1.opus")) + }) + It("zips a playlist correctly", func() { tracks := []model.PlaylistTrack{ {MediaFile: model.MediaFile{Path: "test_data/01 - track1.mp3", Suffix: "mp3", AlbumID: "1", Album: "Album 1", DiscNumber: 1, Artist: "AC/DC", Title: "track1"}}, @@ -571,6 +614,19 @@ func (m *mockMediaStreamer) NewStream(ctx context.Context, mf *model.MediaFile, return &stream.Stream{ReadCloser: args.Get(0).(io.ReadCloser)}, nil } +// fakeDecider echoes the legacy format/bitrate unless a resolved request is set. +type fakeDecider struct { + stream.TranscodeDecider + resolved *stream.Request +} + +func (f *fakeDecider) ResolveRequest(_ context.Context, _ *model.MediaFile, format string, bitRate int, offset int) stream.Request { + if f.resolved != nil { + return *f.resolved + } + return stream.Request{Format: format, BitRate: bitRate, Offset: offset} +} + type mockShare struct { mock.Mock core.Share diff --git a/server/public/handle_streams.go b/server/public/handle_streams.go index 37ae56c2b..46c7ca210 100644 --- a/server/public/handle_streams.go +++ b/server/public/handle_streams.go @@ -60,9 +60,8 @@ func (pub *Router) handleStream(w http.ResponseWriter, r *http.Request) { return } - stream, err := pub.streamer.NewStream(ctx, mf, streampkg.Request{ - Format: info.format, BitRate: info.bitrate, - }) + streamReq := pub.decider.ResolveRequest(ctx, mf, info.format, info.bitrate, 0) + stream, err := pub.streamer.NewStream(ctx, mf, streamReq) if err != nil { if errors.Is(err, streampkg.ErrTooManyTranscodes) { w.Header().Set("Retry-After", strconv.Itoa(streampkg.RetryAfterSeconds)) diff --git a/server/public/handle_streams_test.go b/server/public/handle_streams_test.go index 4b4a3545b..965bc7e05 100644 --- a/server/public/handle_streams_test.go +++ b/server/public/handle_streams_test.go @@ -116,11 +116,11 @@ var _ = Describe("handleStream", func() { BeforeEach(func() { ctx = GinkgoT().Context() auth.PublicTokenAuth = jwtauth.New("HS256", []byte("test-secret"), nil) - ds = &tests.MockDataStore{} + ds = &tests.MockDataStore{MockedTranscoding: &tests.MockTranscodingRepo{}} shareRepo = &tests.MockShareRepo{} ds.MockedShare = shareRepo streamer = &mockStreamer{} - pub = &Router{ds: ds, streamer: streamer} + pub = &Router{ds: ds, streamer: streamer, decider: stream.NewTranscodeDecider(ds, tests.NewMockFFmpeg(""))} }) makeRequest := func(token string) *httptest.ResponseRecorder { @@ -152,8 +152,17 @@ var _ = Describe("handleStream", func() { makeRequest(token) Expect(streamer.called).To(BeTrue()) - Expect(streamer.req.Format).To(Equal("mp3")) - Expect(streamer.req.BitRate).To(Equal(192)) + }) + + It("resolves the full stream request like the Subsonic endpoint, so transcodes share the cache", func() { + mf := model.MediaFile{ID: "mf-123", Suffix: "flac", BitRate: 1500, SampleRate: 44100, BitDepth: new(24), Channels: 2} + shareOwnedBy(model.User{ID: "owner1", UserName: "owner1", IsAdmin: true}, mf) + + claims := auth.Claims{ID: "mf-123", Format: "opus", BitRate: 128, ShareID: "share123"} + token, _ := auth.CreateExpiringPublicToken(time.Now().Add(time.Hour), claims) + makeRequest(token) + + Expect(streamer.req).To(Equal(stream.Request{Format: "opus", BitRate: 128, SampleRate: 48000, Channels: 2})) }) It("returns 404 when the track is outside the share owner's libraries", func() { diff --git a/server/public/public.go b/server/public/public.go index 142c474bd..8239ef927 100644 --- a/server/public/public.go +++ b/server/public/public.go @@ -21,14 +21,15 @@ type Router struct { http.Handler artwork artwork.Artwork streamer stream.MediaStreamer + decider stream.TranscodeDecider archiver core.Archiver share core.Share assetsHandler http.Handler ds model.DataStore } -func New(ds model.DataStore, artwork artwork.Artwork, streamer stream.MediaStreamer, share core.Share, archiver core.Archiver) *Router { - p := &Router{ds: ds, artwork: artwork, streamer: streamer, share: share, archiver: archiver} +func New(ds model.DataStore, artwork artwork.Artwork, streamer stream.MediaStreamer, decider stream.TranscodeDecider, share core.Share, archiver core.Archiver) *Router { + p := &Router{ds: ds, artwork: artwork, streamer: streamer, decider: decider, share: share, archiver: archiver} shareRoot := path.Join(conf.Server.BasePath, consts.URLPathPublic) p.assetsHandler = http.StripPrefix(shareRoot, http.FileServer(http.FS(ui.BuildAssets()))) p.Handler = p.routes() diff --git a/server/subsonic/e2e/subsonic_artwork_test.go b/server/subsonic/e2e/subsonic_artwork_test.go index 9394c8830..530588759 100644 --- a/server/subsonic/e2e/subsonic_artwork_test.go +++ b/server/subsonic/e2e/subsonic_artwork_test.go @@ -137,7 +137,7 @@ var _ = Describe("Artwork Serving", Ordered, func() { artRouter = buildArtworkRouter(artSvc) router = artRouter // so the shared doReq/doRawReq helpers hit the artwork-wired router - pubRouter = public.New(ds, artSvc, streamerSpy, core.NewShare(ds), noopArchiver{}) + pubRouter = public.New(ds, artSvc, streamerSpy, stream.NewTranscodeDecider(ds, ffm), core.NewShare(ds), noopArchiver{}) }) It("emits a bare optimistic coverArt id before the queue is drained", func() {