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() {