From 3259040db06f7d12b28a77405a5f74a5e4cb8614 Mon Sep 17 00:00:00 2001 From: Deluan Date: Sun, 8 Mar 2026 20:06:17 -0400 Subject: [PATCH] refactor(transcode): move MediaStreamer into core/transcode and unify StreamRequest Moved MediaStreamer, Stream, TranscodingCache and related types from core/media_streamer.go into core/transcode/, eliminating the duplicate StreamRequest type. The transcode.StreamRequest now carries all fields (ID, Format, BitRate, SampleRate, BitDepth, Channels, Offset) and ResolveStream returns a fully-populated value, removing manual field copying at every call site. Also moved buildLegacyClientInfo into the transcode package alongside ResolveStream, and unexported ParseTranscodeParams since it was only used internally by ValidateTranscodeParams. --- cmd/wire_gen.go | 10 +-- core/archiver.go | 7 +- core/archiver_test.go | 15 ++-- .../transcode/legacy_client_test.go | 13 ++- core/{ => transcode}/media_streamer.go | 25 +++--- core/{ => transcode}/media_streamer_test.go | 20 ++--- core/transcode/transcode.go | 80 +++++++++++++++++- core/transcode/transcode_test.go | 22 ++--- core/transcode/types.go | 13 ++- core/wire_providers.go | 4 +- server/e2e/e2e_suite_test.go | 26 +++--- server/public/handle_streams.go | 4 +- server/public/public.go | 5 +- server/subsonic/api.go | 4 +- server/subsonic/stream.go | 81 +------------------ server/subsonic/transcode.go | 3 +- server/subsonic/transcode_test.go | 20 ++--- 17 files changed, 180 insertions(+), 172 deletions(-) rename server/subsonic/stream_internal_test.go => core/transcode/legacy_client_test.go (85%) rename core/{ => transcode}/media_streamer.go (94%) rename core/{ => transcode}/media_streamer_test.go (73%) diff --git a/cmd/wire_gen.go b/cmd/wire_gen.go index ebb2031a1..a7a0769d3 100644 --- a/cmd/wire_gen.go +++ b/cmd/wire_gen.go @@ -1,6 +1,6 @@ // Code generated by Wire. DO NOT EDIT. -//go:generate go run -mod=mod github.com/google/wire/cmd/wire gen -tags "netgo" +//go:generate go run -mod=mod github.com/google/wire/cmd/wire gen -tags "netgo sqlite_fts5" //go:build !wireinject // +build !wireinject @@ -95,8 +95,8 @@ func CreateSubsonicAPIRouter(ctx context.Context) *subsonic.Router { agentsAgents := agents.GetAgents(dataStore, manager) provider := external.NewProvider(dataStore, agentsAgents) artworkArtwork := artwork.NewArtwork(dataStore, fileCache, fFmpeg, provider) - transcodingCache := core.GetTranscodingCache() - mediaStreamer := core.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) + transcodingCache := transcode.GetTranscodingCache() + mediaStreamer := transcode.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) share := core.NewShare(dataStore) archiver := core.NewArchiver(mediaStreamer, dataStore, share) players := core.NewPlayers(dataStore) @@ -122,8 +122,8 @@ func CreatePublicRouter() *public.Router { agentsAgents := agents.GetAgents(dataStore, manager) provider := external.NewProvider(dataStore, agentsAgents) artworkArtwork := artwork.NewArtwork(dataStore, fileCache, fFmpeg, provider) - transcodingCache := core.GetTranscodingCache() - mediaStreamer := core.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) + transcodingCache := transcode.GetTranscodingCache() + mediaStreamer := transcode.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) share := core.NewShare(dataStore) archiver := core.NewArchiver(mediaStreamer, dataStore, share) router := public.New(dataStore, artworkArtwork, mediaStreamer, share, archiver) diff --git a/core/archiver.go b/core/archiver.go index fef9188a2..88b2d5b0e 100644 --- a/core/archiver.go +++ b/core/archiver.go @@ -10,6 +10,7 @@ import ( "strings" "github.com/Masterminds/squirrel" + "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/utils/slice" @@ -22,13 +23,13 @@ type Archiver interface { ZipPlaylist(ctx context.Context, id string, format string, bitrate int, w io.Writer) error } -func NewArchiver(ms MediaStreamer, ds model.DataStore, shares Share) Archiver { +func NewArchiver(ms transcode.MediaStreamer, ds model.DataStore, shares Share) Archiver { return &archiver{ds: ds, ms: ms, shares: shares} } type archiver struct { ds model.DataStore - ms MediaStreamer + ms transcode.MediaStreamer shares Share } @@ -176,7 +177,7 @@ func (a *archiver) addFileToZip(ctx context.Context, z *zip.Writer, mf model.Med var r io.ReadCloser if format != "raw" && format != "" { - r, err = a.ms.DoStream(ctx, &mf, StreamRequest{Format: format, BitRate: bitrate}) + r, err = a.ms.DoStream(ctx, &mf, transcode.StreamRequest{Format: format, BitRate: bitrate}) } else { r, err = os.Open(path) } diff --git a/core/archiver_test.go b/core/archiver_test.go index 4ded2eba4..bfce641c9 100644 --- a/core/archiver_test.go +++ b/core/archiver_test.go @@ -9,6 +9,7 @@ import ( "github.com/Masterminds/squirrel" "github.com/navidrome/navidrome/core" + "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/model" . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" @@ -44,7 +45,7 @@ var _ = Describe("Archiver", func() { }}).Return(mfs, nil) ds.On("MediaFile", mock.Anything).Return(mfRepo) - ms.On("DoStream", mock.Anything, mock.Anything, core.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(3) + ms.On("DoStream", mock.Anything, mock.Anything, transcode.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(3) out := new(bytes.Buffer) err := arch.ZipAlbum(context.Background(), "1", "mp3", 128, out) @@ -73,7 +74,7 @@ var _ = Describe("Archiver", func() { }}).Return(mfs, nil) ds.On("MediaFile", mock.Anything).Return(mfRepo) - ms.On("DoStream", mock.Anything, mock.Anything, core.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(2) + ms.On("DoStream", mock.Anything, mock.Anything, transcode.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(2) out := new(bytes.Buffer) err := arch.ZipArtist(context.Background(), "1", "mp3", 128, out) @@ -104,7 +105,7 @@ var _ = Describe("Archiver", func() { } sh.On("Load", mock.Anything, "1").Return(share, nil) - ms.On("DoStream", mock.Anything, mock.Anything, core.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(2) + ms.On("DoStream", mock.Anything, mock.Anything, transcode.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(2) out := new(bytes.Buffer) err := arch.ZipShare(context.Background(), "1", out) @@ -136,7 +137,7 @@ var _ = Describe("Archiver", func() { plRepo := &mockPlaylistRepository{} plRepo.On("GetWithTracks", "1", true, false).Return(pls, nil) ds.On("Playlist", mock.Anything).Return(plRepo) - ms.On("DoStream", mock.Anything, mock.Anything, core.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(2) + ms.On("DoStream", mock.Anything, mock.Anything, transcode.StreamRequest{Format: "mp3", BitRate: 128}).Return(io.NopCloser(strings.NewReader("test")), nil).Times(2) out := new(bytes.Buffer) err := arch.ZipPlaylist(context.Background(), "1", "mp3", 128, out) @@ -214,15 +215,15 @@ func (m *mockPlaylistRepository) GetWithTracks(id string, refreshSmartPlaylists, type mockMediaStreamer struct { mock.Mock - core.MediaStreamer + transcode.MediaStreamer } -func (m *mockMediaStreamer) DoStream(ctx context.Context, mf *model.MediaFile, req core.StreamRequest) (*core.Stream, error) { +func (m *mockMediaStreamer) DoStream(ctx context.Context, mf *model.MediaFile, req transcode.StreamRequest) (*transcode.Stream, error) { args := m.Called(ctx, mf, req) if args.Error(1) != nil { return nil, args.Error(1) } - return &core.Stream{ReadCloser: args.Get(0).(io.ReadCloser)}, nil + return &transcode.Stream{ReadCloser: args.Get(0).(io.ReadCloser)}, nil } type mockShare struct { diff --git a/server/subsonic/stream_internal_test.go b/core/transcode/legacy_client_test.go similarity index 85% rename from server/subsonic/stream_internal_test.go rename to core/transcode/legacy_client_test.go index aeb8c90e7..9628764f4 100644 --- a/server/subsonic/stream_internal_test.go +++ b/core/transcode/legacy_client_test.go @@ -1,9 +1,8 @@ -package subsonic +package transcode import ( "github.com/navidrome/navidrome/conf" "github.com/navidrome/navidrome/conf/configtest" - "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/model" . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" @@ -23,13 +22,13 @@ var _ = Describe("buildLegacyClientInfo", func() { Expect(ci.TranscodingProfiles).To(HaveLen(1)) Expect(ci.TranscodingProfiles[0].Container).To(Equal("mp3")) Expect(ci.TranscodingProfiles[0].AudioCodec).To(Equal("mp3")) - Expect(ci.TranscodingProfiles[0].Protocol).To(Equal(transcode.ProtocolHTTP)) + Expect(ci.TranscodingProfiles[0].Protocol).To(Equal(ProtocolHTTP)) Expect(ci.MaxAudioBitrate).To(BeZero()) Expect(ci.MaxTranscodingAudioBitrate).To(BeZero()) Expect(ci.DirectPlayProfiles).To(HaveLen(1)) Expect(ci.DirectPlayProfiles[0].Containers).To(Equal([]string{"flac"})) Expect(ci.DirectPlayProfiles[0].AudioCodecs).To(Equal([]string{mf.AudioCodec()})) - Expect(ci.DirectPlayProfiles[0].Protocols).To(Equal([]string{transcode.ProtocolHTTP})) + Expect(ci.DirectPlayProfiles[0].Protocols).To(Equal([]string{ProtocolHTTP})) }) It("sets transcoding profile and bitrate for explicit format with bitrate", func() { @@ -50,7 +49,7 @@ var _ = Describe("buildLegacyClientInfo", func() { Expect(ci.DirectPlayProfiles).To(HaveLen(1)) Expect(ci.DirectPlayProfiles[0].Containers).To(BeEmpty()) Expect(ci.DirectPlayProfiles[0].AudioCodecs).To(BeEmpty()) - Expect(ci.DirectPlayProfiles[0].Protocols).To(Equal([]string{transcode.ProtocolHTTP})) + Expect(ci.DirectPlayProfiles[0].Protocols).To(Equal([]string{ProtocolHTTP})) Expect(ci.TranscodingProfiles).To(BeEmpty()) Expect(ci.MaxAudioBitrate).To(BeZero()) }) @@ -64,7 +63,7 @@ var _ = Describe("buildLegacyClientInfo", func() { Expect(ci.TranscodingProfiles).To(HaveLen(1)) Expect(ci.TranscodingProfiles[0].Container).To(Equal("opus")) Expect(ci.TranscodingProfiles[0].AudioCodec).To(Equal("opus")) - Expect(ci.TranscodingProfiles[0].Protocol).To(Equal(transcode.ProtocolHTTP)) + Expect(ci.TranscodingProfiles[0].Protocol).To(Equal(ProtocolHTTP)) Expect(ci.MaxAudioBitrate).To(Equal(128)) Expect(ci.MaxTranscodingAudioBitrate).To(Equal(128)) Expect(ci.DirectPlayProfiles).To(HaveLen(1)) @@ -78,7 +77,7 @@ var _ = Describe("buildLegacyClientInfo", func() { Expect(ci.DirectPlayProfiles).To(HaveLen(1)) Expect(ci.DirectPlayProfiles[0].Containers).To(BeEmpty()) Expect(ci.DirectPlayProfiles[0].AudioCodecs).To(BeEmpty()) - Expect(ci.DirectPlayProfiles[0].Protocols).To(Equal([]string{transcode.ProtocolHTTP})) + Expect(ci.DirectPlayProfiles[0].Protocols).To(Equal([]string{ProtocolHTTP})) Expect(ci.TranscodingProfiles).To(BeEmpty()) Expect(ci.MaxAudioBitrate).To(BeZero()) }) diff --git a/core/media_streamer.go b/core/transcode/media_streamer.go similarity index 94% rename from core/media_streamer.go rename to core/transcode/media_streamer.go index 59f91821e..3f50d43ba 100644 --- a/core/media_streamer.go +++ b/core/transcode/media_streamer.go @@ -1,4 +1,4 @@ -package core +package transcode import ( "context" @@ -12,24 +12,12 @@ import ( "github.com/navidrome/navidrome/conf" "github.com/navidrome/navidrome/consts" "github.com/navidrome/navidrome/core/ffmpeg" - "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/model/request" "github.com/navidrome/navidrome/utils/cache" ) -// StreamRequest contains all parameters for creating a media stream. -type StreamRequest struct { - ID string - Format string - BitRate int // kbps - SampleRate int - BitDepth int - Channels int - Offset int // seconds -} - type MediaStreamer interface { NewStream(ctx context.Context, req StreamRequest) (*Stream, error) DoStream(ctx context.Context, mf *model.MediaFile, req StreamRequest) (*Stream, error) @@ -171,7 +159,7 @@ func NewTranscodingCache() TranscodingCache { consts.TranscodingCacheDir, consts.DefaultTranscodingCacheMaxItems, func(ctx context.Context, arg cache.Item) (io.Reader, error) { job := arg.(*streamJob) - command := transcode.LookupTranscodeCommand(ctx, job.ms.ds, job.format) + command := LookupTranscodeCommand(ctx, job.ms.ds, job.format) if command == "" { log.Error(ctx, "No transcoding command available", "format", job.format) return nil, os.ErrInvalid @@ -206,3 +194,12 @@ func NewTranscodingCache() TranscodingCache { return out, nil }) } + +// userName extracts the username from the context for logging purposes. +func userName(ctx context.Context) string { + if user, ok := request.UserFrom(ctx); !ok { + return "UNKNOWN" + } else { + return user.UserName + } +} diff --git a/core/media_streamer_test.go b/core/transcode/media_streamer_test.go similarity index 73% rename from core/media_streamer_test.go rename to core/transcode/media_streamer_test.go index e47beb66d..f49dcb8d8 100644 --- a/core/media_streamer_test.go +++ b/core/transcode/media_streamer_test.go @@ -1,4 +1,4 @@ -package core_test +package transcode_test import ( "context" @@ -7,7 +7,7 @@ import ( "github.com/navidrome/navidrome/conf" "github.com/navidrome/navidrome/conf/configtest" - "github.com/navidrome/navidrome/core" + "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/tests" @@ -16,7 +16,7 @@ import ( ) var _ = Describe("MediaStreamer", func() { - var streamer core.MediaStreamer + var streamer transcode.MediaStreamer var ds model.DataStore ffmpeg := tests.NewMockFFmpeg("fake data") ctx := log.NewContext(context.TODO()) @@ -29,9 +29,9 @@ var _ = Describe("MediaStreamer", func() { ds.MediaFile(ctx).(*tests.MockMediaFileRepo).SetData(model.MediaFiles{ {ID: "123", Path: "tests/fixtures/test.mp3", Suffix: "mp3", BitRate: 128, Duration: 257.0}, }) - testCache := core.NewTranscodingCache() + testCache := transcode.NewTranscodingCache() Eventually(func() bool { return testCache.Available(context.TODO()) }).Should(BeTrue()) - streamer = core.NewMediaStreamer(ds, ffmpeg, testCache) + streamer = transcode.NewMediaStreamer(ds, ffmpeg, testCache) }) AfterEach(func() { _ = os.RemoveAll(conf.Server.CacheFolder) @@ -39,29 +39,29 @@ var _ = Describe("MediaStreamer", func() { Context("NewStream", func() { It("returns a seekable stream if format is 'raw'", func() { - s, err := streamer.NewStream(ctx, core.StreamRequest{ID: "123", Format: "raw"}) + s, err := streamer.NewStream(ctx, transcode.StreamRequest{ID: "123", Format: "raw"}) Expect(err).ToNot(HaveOccurred()) Expect(s.Seekable()).To(BeTrue()) }) It("returns a seekable stream if no format is specified (direct play)", func() { - s, err := streamer.NewStream(ctx, core.StreamRequest{ID: "123"}) + s, err := streamer.NewStream(ctx, transcode.StreamRequest{ID: "123"}) Expect(err).ToNot(HaveOccurred()) Expect(s.Seekable()).To(BeTrue()) }) It("returns a NON seekable stream if transcode is required", func() { - s, err := streamer.NewStream(ctx, core.StreamRequest{ID: "123", Format: "mp3", BitRate: 64}) + s, err := streamer.NewStream(ctx, transcode.StreamRequest{ID: "123", Format: "mp3", BitRate: 64}) Expect(err).To(BeNil()) Expect(s.Seekable()).To(BeFalse()) Expect(s.Duration()).To(Equal(float32(257.0))) }) It("returns a seekable stream if the file is complete in the cache", func() { - s, err := streamer.NewStream(ctx, core.StreamRequest{ID: "123", Format: "mp3", BitRate: 32}) + s, err := streamer.NewStream(ctx, transcode.StreamRequest{ID: "123", Format: "mp3", BitRate: 32}) Expect(err).To(BeNil()) _, _ = io.ReadAll(s) _ = s.Close() Eventually(func() bool { return ffmpeg.IsClosed() }, "3s").Should(BeTrue()) - s, err = streamer.NewStream(ctx, core.StreamRequest{ID: "123", Format: "mp3", BitRate: 32}) + s, err = streamer.NewStream(ctx, transcode.StreamRequest{ID: "123", Format: "mp3", BitRate: 32}) Expect(err).To(BeNil()) Expect(s.Seekable()).To(BeTrue()) }) diff --git a/core/transcode/transcode.go b/core/transcode/transcode.go index 1fd48f8ca..0039333b2 100644 --- a/core/transcode/transcode.go +++ b/core/transcode/transcode.go @@ -423,11 +423,87 @@ func (s *deciderService) ensureProbed(ctx context.Context, mf *model.MediaFile) return result, nil } +// buildLegacyClientInfo translates legacy Subsonic stream/download parameters +// into a ClientInfo for use with MakeDecision. +// It does NOT read request.TranscodingFrom(ctx) — that is handled by +// MakeDecision's applyServerOverride. +func buildLegacyClientInfo(mf *model.MediaFile, reqFormat string, reqBitRate int) *ClientInfo { + ci := &ClientInfo{Name: "legacy"} + + // Determine target format for transcoding + var targetFormat string + switch { + case reqFormat != "": + targetFormat = reqFormat + case reqBitRate > 0 && reqBitRate < mf.BitRate && conf.Server.DefaultDownsamplingFormat != "": + targetFormat = conf.Server.DefaultDownsamplingFormat + } + + if targetFormat != "" { + ci.DirectPlayProfiles = []DirectPlayProfile{ + {Containers: []string{mf.Suffix}, AudioCodecs: []string{mf.AudioCodec()}, Protocols: []string{ProtocolHTTP}}, + } + ci.TranscodingProfiles = []Profile{ + {Container: targetFormat, AudioCodec: targetFormat, Protocol: ProtocolHTTP}, + } + if reqBitRate > 0 { + ci.MaxAudioBitrate = reqBitRate + ci.MaxTranscodingAudioBitrate = reqBitRate + } + } else { + // No transcoding requested — direct play everything + ci.DirectPlayProfiles = []DirectPlayProfile{ + {Protocols: []string{ProtocolHTTP}}, + } + } + + return ci +} + +// ResolveStream uses MakeDecision to resolve legacy Subsonic stream parameters +// into a fully specified StreamRequest. +func (s *deciderService) ResolveStream(ctx context.Context, mf *model.MediaFile, reqFormat string, reqBitRate int, offset int) StreamRequest { + var req StreamRequest + req.ID = mf.ID + req.Offset = offset + + if reqFormat == "raw" { + req.Format = "raw" + return req + } + + clientInfo := buildLegacyClientInfo(mf, reqFormat, reqBitRate) + decision, err := s.MakeDecision(ctx, mf, clientInfo, DecisionOptions{SkipProbe: true}) + if err != nil { + log.Error(ctx, "Error making transcode decision, falling back to raw", "id", mf.ID, err) + req.Format = "raw" + return req + } + + if decision.CanDirectPlay { + req.Format = "raw" + return req + } + + if decision.CanTranscode { + req.Format = decision.TargetFormat + req.BitRate = decision.TargetBitrate + req.SampleRate = decision.TargetSampleRate + req.BitDepth = decision.TargetBitDepth + req.Channels = decision.TargetChannels + return req + } + + // No compatible profile — fallback to raw + req.Format = "raw" + return req +} + func (s *deciderService) CreateTranscodeParams(decision *Decision) (string, error) { return auth.EncodeToken(decision.toClaimsMap()) } -func (s *deciderService) ParseTranscodeParams(tokenStr string) (*Params, error) { +func (s *deciderService) parseTranscodeParams(tokenStr string) (*Params, error) { token, err := auth.DecodeAndVerifyToken(tokenStr) if err != nil { return nil, err @@ -436,7 +512,7 @@ func (s *deciderService) ParseTranscodeParams(tokenStr string) (*Params, error) } func (s *deciderService) ValidateTranscodeParams(ctx context.Context, token string, mediaID string) (*Params, *model.MediaFile, error) { - params, err := s.ParseTranscodeParams(token) + params, err := s.parseTranscodeParams(token) if err != nil { return nil, nil, errors.Join(ErrTokenInvalid, err) } diff --git a/core/transcode/transcode_test.go b/core/transcode/transcode_test.go index 35e8f18bf..699fe45f4 100644 --- a/core/transcode/transcode_test.go +++ b/core/transcode/transcode_test.go @@ -1101,10 +1101,14 @@ var _ = Describe("Decider", func() { }) Describe("Token round-trip", func() { - var sourceTime time.Time + var ( + sourceTime time.Time + impl *deciderService + ) BeforeEach(func() { sourceTime = time.Date(2025, 6, 15, 10, 30, 0, 0, time.UTC) + impl = svc.(*deciderService) }) It("creates and parses a direct play token", func() { @@ -1117,7 +1121,7 @@ var _ = Describe("Decider", func() { Expect(err).ToNot(HaveOccurred()) Expect(token).ToNot(BeEmpty()) - params, err := svc.ParseTranscodeParams(token) + params, err := impl.parseTranscodeParams(token) Expect(err).ToNot(HaveOccurred()) Expect(params.MediaID).To(Equal("media-123")) Expect(params.DirectPlay).To(BeTrue()) @@ -1138,7 +1142,7 @@ var _ = Describe("Decider", func() { token, err := svc.CreateTranscodeParams(decision) Expect(err).ToNot(HaveOccurred()) - params, err := svc.ParseTranscodeParams(token) + params, err := impl.parseTranscodeParams(token) Expect(err).ToNot(HaveOccurred()) Expect(params.MediaID).To(Equal("media-456")) Expect(params.DirectPlay).To(BeFalse()) @@ -1162,7 +1166,7 @@ var _ = Describe("Decider", func() { token, err := svc.CreateTranscodeParams(decision) Expect(err).ToNot(HaveOccurred()) - params, err := svc.ParseTranscodeParams(token) + params, err := impl.parseTranscodeParams(token) Expect(err).ToNot(HaveOccurred()) Expect(params.MediaID).To(Equal("media-789")) Expect(params.DirectPlay).To(BeFalse()) @@ -1185,7 +1189,7 @@ var _ = Describe("Decider", func() { token, err := svc.CreateTranscodeParams(decision) Expect(err).ToNot(HaveOccurred()) - params, err := svc.ParseTranscodeParams(token) + params, err := impl.parseTranscodeParams(token) Expect(err).ToNot(HaveOccurred()) Expect(params.MediaID).To(Equal("media-bd")) Expect(params.TargetBitDepth).To(Equal(24)) @@ -1204,7 +1208,7 @@ var _ = Describe("Decider", func() { token, err := svc.CreateTranscodeParams(decision) Expect(err).ToNot(HaveOccurred()) - params, err := svc.ParseTranscodeParams(token) + params, err := impl.parseTranscodeParams(token) Expect(err).ToNot(HaveOccurred()) Expect(params.TargetBitDepth).To(Equal(0)) }) @@ -1222,7 +1226,7 @@ var _ = Describe("Decider", func() { token, err := svc.CreateTranscodeParams(decision) Expect(err).ToNot(HaveOccurred()) - params, err := svc.ParseTranscodeParams(token) + params, err := impl.parseTranscodeParams(token) Expect(err).ToNot(HaveOccurred()) Expect(params.TargetSampleRate).To(Equal(0)) }) @@ -1237,13 +1241,13 @@ var _ = Describe("Decider", func() { token, err := svc.CreateTranscodeParams(decision) Expect(err).ToNot(HaveOccurred()) - params, err := svc.ParseTranscodeParams(token) + params, err := impl.parseTranscodeParams(token) Expect(err).ToNot(HaveOccurred()) Expect(params.SourceUpdatedAt.Unix()).To(Equal(timeWithNanos.Truncate(time.Second).Unix())) }) It("rejects an invalid token", func() { - _, err := svc.ParseTranscodeParams("invalid-token") + _, err := impl.parseTranscodeParams("invalid-token") Expect(err).To(HaveOccurred()) }) }) diff --git a/core/transcode/types.go b/core/transcode/types.go index ffc79349e..4e24a17c9 100644 --- a/core/transcode/types.go +++ b/core/transcode/types.go @@ -23,11 +23,22 @@ type DecisionOptions struct { SkipProbe bool } +// StreamRequest contains the resolved parameters for creating a media stream. +type StreamRequest struct { + ID string + Format string + BitRate int // kbps + SampleRate int + BitDepth int + Channels int + Offset int // seconds +} + // Decider is the core service interface for making transcoding decisions type Decider interface { MakeDecision(ctx context.Context, mf *model.MediaFile, clientInfo *ClientInfo, opts DecisionOptions) (*Decision, error) + ResolveStream(ctx context.Context, mf *model.MediaFile, reqFormat string, reqBitRate int, offset int) StreamRequest CreateTranscodeParams(decision *Decision) (string, error) - ParseTranscodeParams(token string) (*Params, error) ValidateTranscodeParams(ctx context.Context, token string, mediaID string) (*Params, *model.MediaFile, error) } diff --git a/core/wire_providers.go b/core/wire_providers.go index 7ba7c879b..20b5eb9a5 100644 --- a/core/wire_providers.go +++ b/core/wire_providers.go @@ -14,8 +14,8 @@ import ( ) var Set = wire.NewSet( - NewMediaStreamer, - GetTranscodingCache, + transcode.NewMediaStreamer, + transcode.GetTranscodingCache, NewArchiver, NewPlayers, NewShare, diff --git a/server/e2e/e2e_suite_test.go b/server/e2e/e2e_suite_test.go index 041db0935..4be79f940 100644 --- a/server/e2e/e2e_suite_test.go +++ b/server/e2e/e2e_suite_test.go @@ -225,14 +225,14 @@ func (n noopArtwork) GetOrPlaceholder(_ context.Context, _ string, _ int, _ bool return io.NopCloser(io.LimitReader(nil, 0)), time.Time{}, nil } -// noopStreamer implements core.MediaStreamer +// noopStreamer implements transcode.MediaStreamer type noopStreamer struct{} -func (n noopStreamer) NewStream(context.Context, core.StreamRequest) (*core.Stream, error) { +func (n noopStreamer) NewStream(context.Context, transcode.StreamRequest) (*transcode.Stream, error) { return nil, model.ErrNotFound } -func (n noopStreamer) DoStream(context.Context, *model.MediaFile, core.StreamRequest) (*core.Stream, error) { +func (n noopStreamer) DoStream(context.Context, *model.MediaFile, transcode.StreamRequest) (*transcode.Stream, error) { return nil, model.ErrNotFound } @@ -243,12 +243,12 @@ func (n noopDecider) MakeDecision(context.Context, *model.MediaFile, *transcode. return nil, nil } -func (n noopDecider) CreateTranscodeParams(*transcode.Decision) (string, error) { - return "", nil +func (n noopDecider) ResolveStream(context.Context, *model.MediaFile, string, int, int) transcode.StreamRequest { + return transcode.StreamRequest{Format: "raw"} } -func (n noopDecider) ParseTranscodeParams(string) (*transcode.Params, error) { - return nil, nil +func (n noopDecider) CreateTranscodeParams(*transcode.Decision) (string, error) { + return "", nil } func (n noopDecider) ValidateTranscodeParams(context.Context, string, string) (*transcode.Params, *model.MediaFile, error) { @@ -318,12 +318,12 @@ func (n noopPlayTracker) Submit(context.Context, []scrobbler.Submission) error { // Compile-time interface checks var ( - _ artwork.Artwork = noopArtwork{} - _ core.MediaStreamer = noopStreamer{} - _ core.Archiver = noopArchiver{} - _ external.Provider = noopProvider{} - _ scrobbler.PlayTracker = noopPlayTracker{} - _ transcode.Decider = noopDecider{} + _ artwork.Artwork = noopArtwork{} + _ transcode.MediaStreamer = noopStreamer{} + _ core.Archiver = noopArchiver{} + _ external.Provider = noopProvider{} + _ scrobbler.PlayTracker = noopPlayTracker{} + _ transcode.Decider = noopDecider{} ) var _ = BeforeSuite(func() { diff --git a/server/public/handle_streams.go b/server/public/handle_streams.go index 2baec52ed..a147a2ac8 100644 --- a/server/public/handle_streams.go +++ b/server/public/handle_streams.go @@ -6,8 +6,8 @@ import ( "net/http" "strconv" - "github.com/navidrome/navidrome/core" "github.com/navidrome/navidrome/core/auth" + "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/utils/req" ) @@ -23,7 +23,7 @@ func (pub *Router) handleStream(w http.ResponseWriter, r *http.Request) { return } - stream, err := pub.streamer.NewStream(ctx, core.StreamRequest{ + stream, err := pub.streamer.NewStream(ctx, transcode.StreamRequest{ ID: info.id, Format: info.format, BitRate: info.bitrate, }) if err != nil { diff --git a/server/public/public.go b/server/public/public.go index ebccb01d2..7d8a4e007 100644 --- a/server/public/public.go +++ b/server/public/public.go @@ -11,6 +11,7 @@ import ( "github.com/navidrome/navidrome/core" "github.com/navidrome/navidrome/core/artwork" "github.com/navidrome/navidrome/core/publicurl" + "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/server" @@ -20,14 +21,14 @@ import ( type Router struct { http.Handler artwork artwork.Artwork - streamer core.MediaStreamer + streamer transcode.MediaStreamer archiver core.Archiver share core.Share assetsHandler http.Handler ds model.DataStore } -func New(ds model.DataStore, artwork artwork.Artwork, streamer core.MediaStreamer, share core.Share, archiver core.Archiver) *Router { +func New(ds model.DataStore, artwork artwork.Artwork, streamer transcode.MediaStreamer, share core.Share, archiver core.Archiver) *Router { p := &Router{ds: ds, artwork: artwork, streamer: streamer, share: share, archiver: archiver} shareRoot := path.Join(conf.Server.BasePath, consts.URLPathPublic) p.assetsHandler = http.StripPrefix(shareRoot, http.FileServer(http.FS(ui.BuildAssets()))) diff --git a/server/subsonic/api.go b/server/subsonic/api.go index 584f752b5..6f355d161 100644 --- a/server/subsonic/api.go +++ b/server/subsonic/api.go @@ -39,7 +39,7 @@ type Router struct { http.Handler ds model.DataStore artwork artwork.Artwork - streamer core.MediaStreamer + streamer transcode.MediaStreamer archiver core.Archiver players core.Players provider external.Provider @@ -54,7 +54,7 @@ type Router struct { transcodeDecision transcode.Decider } -func New(ds model.DataStore, artwork artwork.Artwork, streamer core.MediaStreamer, archiver core.Archiver, +func New(ds model.DataStore, artwork artwork.Artwork, streamer transcode.MediaStreamer, archiver core.Archiver, players core.Players, provider external.Provider, scanner model.Scanner, broker events.Broker, playlists playlistsvc.Playlists, scrobbler scrobbler.PlayTracker, share core.Share, playback playback.PlaybackServer, metrics metrics.Metrics, lyrics lyricssvc.Lyrics, transcodeDecision transcode.Decider, diff --git a/server/subsonic/stream.go b/server/subsonic/stream.go index 38824ee69..78ccb28d7 100644 --- a/server/subsonic/stream.go +++ b/server/subsonic/stream.go @@ -9,7 +9,6 @@ import ( "strings" "github.com/navidrome/navidrome/conf" - "github.com/navidrome/navidrome/core" "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" @@ -18,7 +17,7 @@ import ( "github.com/navidrome/navidrome/utils/req" ) -func (api *Router) serveStream(ctx context.Context, w http.ResponseWriter, r *http.Request, stream *core.Stream, id string) { +func (api *Router) serveStream(ctx context.Context, w http.ResponseWriter, r *http.Request, stream *transcode.Stream, id string) { if stream.Seekable() { http.ServeContent(w, r, stream.Name(), stream.ModTime(), stream) } else { @@ -66,7 +65,7 @@ func (api *Router) Stream(w http.ResponseWriter, r *http.Request) (*responses.Su return nil, err } - streamReq := api.resolveStreamRequest(ctx, mf, format, maxBitRate, timeOffset) + streamReq := api.transcodeDecision.ResolveStream(ctx, mf, format, maxBitRate, timeOffset) stream, err := api.streamer.DoStream(ctx, mf, streamReq) if err != nil { return nil, err @@ -136,7 +135,7 @@ func (api *Router) Download(w http.ResponseWriter, r *http.Request) (*responses. switch v := entity.(type) { case *model.MediaFile: - streamReq := api.resolveStreamRequest(ctx, v, format, maxBitRate, 0) + streamReq := api.transcodeDecision.ResolveStream(ctx, v, format, maxBitRate, 0) stream, err := api.streamer.DoStream(ctx, v, streamReq) if err != nil { return nil, err @@ -169,77 +168,3 @@ func (api *Router) Download(w http.ResponseWriter, r *http.Request) (*responses. return nil, err } - -// buildLegacyClientInfo translates legacy Subsonic stream/download parameters -// into a transcode.ClientInfo for use with MakeDecision. -// It does NOT read request.TranscodingFrom(ctx) — that is handled by -// MakeDecision's applyServerOverride. -func buildLegacyClientInfo(mf *model.MediaFile, reqFormat string, reqBitRate int) *transcode.ClientInfo { - ci := &transcode.ClientInfo{Name: "legacy"} - - // Determine target format for transcoding - var targetFormat string - switch { - case reqFormat != "": - targetFormat = reqFormat - case reqBitRate > 0 && reqBitRate < mf.BitRate && conf.Server.DefaultDownsamplingFormat != "": - targetFormat = conf.Server.DefaultDownsamplingFormat - } - - if targetFormat != "" { - ci.DirectPlayProfiles = []transcode.DirectPlayProfile{ - {Containers: []string{mf.Suffix}, AudioCodecs: []string{mf.AudioCodec()}, Protocols: []string{transcode.ProtocolHTTP}}, - } - ci.TranscodingProfiles = []transcode.Profile{ - {Container: targetFormat, AudioCodec: targetFormat, Protocol: transcode.ProtocolHTTP}, - } - if reqBitRate > 0 { - ci.MaxAudioBitrate = reqBitRate - ci.MaxTranscodingAudioBitrate = reqBitRate - } - } else { - // No transcoding requested — direct play everything - ci.DirectPlayProfiles = []transcode.DirectPlayProfile{ - {Protocols: []string{transcode.ProtocolHTTP}}, - } - } - - return ci -} - -// resolveStreamRequest uses MakeDecision to resolve legacy stream parameters -// into a fully specified StreamRequest. -func (api *Router) resolveStreamRequest(ctx context.Context, mf *model.MediaFile, reqFormat string, reqBitRate int, offset int) core.StreamRequest { - req := core.StreamRequest{ID: mf.ID, Offset: offset} - - if reqFormat == "raw" { - req.Format = "raw" - return req - } - - clientInfo := buildLegacyClientInfo(mf, reqFormat, reqBitRate) - decision, err := api.transcodeDecision.MakeDecision(ctx, mf, clientInfo, transcode.DecisionOptions{SkipProbe: true}) - if err != nil { - log.Error(ctx, "Error making transcode decision, falling back to raw", "id", mf.ID, err) - req.Format = "raw" - return req - } - - if decision.CanDirectPlay { - req.Format = "raw" - return req - } - - if decision.CanTranscode { - req.Format = decision.TargetFormat - req.BitRate = decision.TargetBitrate - req.SampleRate = decision.TargetSampleRate - req.BitDepth = decision.TargetBitDepth - req.Channels = decision.TargetChannels - return req - } - - // No compatible profile — fallback to raw - req.Format = "raw" - return req -} diff --git a/server/subsonic/transcode.go b/server/subsonic/transcode.go index 1c8d52aea..429b859a5 100644 --- a/server/subsonic/transcode.go +++ b/server/subsonic/transcode.go @@ -8,7 +8,6 @@ import ( "slices" "strconv" - "github.com/navidrome/navidrome/core" "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" @@ -360,7 +359,7 @@ func (api *Router) GetTranscodeStream(w http.ResponseWriter, r *http.Request) (* } // Build streaming parameters from the token - streamReq := core.StreamRequest{ID: mediaID, Offset: p.IntOr("offset", 0)} + streamReq := transcode.StreamRequest{ID: mediaID, Offset: p.IntOr("offset", 0)} if !params.DirectPlay && params.TargetFormat != "" { streamReq.Format = params.TargetFormat streamReq.BitRate = params.TargetBitrate // Already in kbps, matching the streamer diff --git a/server/subsonic/transcode_test.go b/server/subsonic/transcode_test.go index ca86f0a78..062837974 100644 --- a/server/subsonic/transcode_test.go +++ b/server/subsonic/transcode_test.go @@ -7,7 +7,6 @@ import ( "net/http" "net/http/httptest" - "github.com/navidrome/navidrome/core" "github.com/navidrome/navidrome/core/transcode" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/tests" @@ -369,8 +368,6 @@ type mockTranscodeDecision struct { decision *transcode.Decision token string tokenErr error - params *transcode.Params - parseErr error validateParams *transcode.Params validateMF *model.MediaFile validateErr error @@ -383,15 +380,12 @@ func (m *mockTranscodeDecision) MakeDecision(_ context.Context, _ *model.MediaFi return &transcode.Decision{}, nil } -func (m *mockTranscodeDecision) CreateTranscodeParams(_ *transcode.Decision) (string, error) { - return m.token, m.tokenErr +func (m *mockTranscodeDecision) ResolveStream(_ context.Context, _ *model.MediaFile, _ string, _ int, _ int) transcode.StreamRequest { + return transcode.StreamRequest{Format: "raw"} } -func (m *mockTranscodeDecision) ParseTranscodeParams(_ string) (*transcode.Params, error) { - if m.parseErr != nil { - return nil, m.parseErr - } - return m.params, nil +func (m *mockTranscodeDecision) CreateTranscodeParams(_ *transcode.Decision) (string, error) { + return m.token, m.tokenErr } func (m *mockTranscodeDecision) ValidateTranscodeParams(_ context.Context, _ string, _ string) (*transcode.Params, *model.MediaFile, error) { @@ -406,15 +400,15 @@ func (m *mockTranscodeDecision) ValidateTranscodeParams(_ context.Context, _ str var errStreamCaptured = errors.New("stream request captured") type fakeMediaStreamer struct { - captured *core.StreamRequest + captured *transcode.StreamRequest } -func (f *fakeMediaStreamer) NewStream(_ context.Context, req core.StreamRequest) (*core.Stream, error) { +func (f *fakeMediaStreamer) NewStream(_ context.Context, req transcode.StreamRequest) (*transcode.Stream, error) { f.captured = &req return nil, errStreamCaptured } -func (f *fakeMediaStreamer) DoStream(_ context.Context, _ *model.MediaFile, req core.StreamRequest) (*core.Stream, error) { +func (f *fakeMediaStreamer) DoStream(_ context.Context, _ *model.MediaFile, req transcode.StreamRequest) (*transcode.Stream, error) { f.captured = &req return nil, errStreamCaptured }