From 8dee15501fa8ec59a64b3293cfaf3ffd8b64d180 Mon Sep 17 00:00:00 2001 From: Deluan Date: Thu, 24 Sep 2026 16:50:41 -0400 Subject: [PATCH] refactor(transcoding): repair piped headers behind one patcher The FLAC and mp3 repairs had grown into two copies of the same wrapper: the same lazy peek on the first Read, the same MultiReader hand-off, the same tolerance of a short read, and the same comment twice. They now share one headerPatcher, with a prefix function per container. Transcode no longer selects the patcher by opts.Format. That string is the target format of a transcoding profile, which users can name anything, so a row called 'mp3 320' running the default mp3 command got no repair at all, while a row whose command emits FLAC got one only because the FLAC arm was the default. The shared peek dispatches on the magic bytes instead, which both patchers already had to check, and leaves anything it does not recognize untouched. Also drops b2i in favour of gg.If, and lowers the ID3 bound to 64KB: ffmpeg cannot write an attached picture to a non-seekable output, so the tag ffmpeg pipes out stays small. --- core/ffmpeg/ffmpeg.go | 10 +--- core/ffmpeg/flac_streaminfo.go | 43 ++++------------- core/ffmpeg/flac_streaminfo_test.go | 6 +-- core/ffmpeg/mp3_xing.go | 74 +++++------------------------ core/ffmpeg/mp3_xing_test.go | 30 ++++++------ core/ffmpeg/piped_header.go | 61 ++++++++++++++++++++++++ 6 files changed, 102 insertions(+), 122 deletions(-) create mode 100644 core/ffmpeg/piped_header.go diff --git a/core/ffmpeg/ffmpeg.go b/core/ffmpeg/ffmpeg.go index c56ad1a69..fbc66688b 100644 --- a/core/ffmpeg/ffmpeg.go +++ b/core/ffmpeg/ffmpeg.go @@ -32,7 +32,7 @@ type TranscodeOptions struct { Channels int // 0 = no constraint BitDepth int // 0 = no constraint; valid values: 16, 24, 32 Offset int // seconds - Duration float32 // seconds; 0 = unknown. Only used to repair a piped FLAC or mp3 header. + Duration float32 // seconds; 0 = unknown. Only used to repair a header ffmpeg piped out unfinished. } // AudioProbeResult contains authoritative audio stream properties from ffprobe. @@ -91,13 +91,7 @@ func (e *ffmpeg) Transcode(ctx context.Context, opts TranscodeOptions) (io.ReadC if err != nil { return nil, err } - duration := opts.Duration - float32(opts.Offset) - switch opts.Format { - case "mp3": - return patchMP3Duration(out, duration), nil - default: - return patchFLACDuration(out, duration), nil - } + return patchPipedHeader(out, opts.Duration-float32(opts.Offset)), nil } func (e *ffmpeg) ConvertAnimatedImage(ctx context.Context, reader io.Reader, maxSize int, quality int) (io.ReadCloser, error) { diff --git a/core/ffmpeg/flac_streaminfo.go b/core/ffmpeg/flac_streaminfo.go index 878c28718..f4e70f3a9 100644 --- a/core/ffmpeg/flac_streaminfo.go +++ b/core/ffmpeg/flac_streaminfo.go @@ -1,10 +1,7 @@ package ffmpeg import ( - "bytes" "encoding/binary" - "errors" - "io" "math" ) @@ -13,43 +10,21 @@ const ( flacMaxTotalSamples = 1<<36 - 1 ) -// patchFLACDuration fills in the STREAMINFO total_samples that ffmpeg leaves at 0 -// when writing to a pipe, since a decoder cannot seek a cached FLAC without it. -func patchFLACDuration(r io.ReadCloser, duration float32) io.ReadCloser { - if duration <= 0 { - return r +// flacPrefix fills in the STREAMINFO total_samples that ffmpeg leaves at 0 when +// writing to a pipe, since a decoder cannot seek a cached FLAC without it. +func (h *headerPatcher) flacPrefix(buf []byte) ([]byte, error) { + buf, err := h.fill(buf, flacPrefixLen) + if err != nil || len(buf) < flacPrefixLen { + return buf, err } - return &flacPatcher{ReadCloser: r, duration: duration} -} - -type flacPatcher struct { - io.ReadCloser - duration float32 - // Peeking here rather than in the constructor keeps Transcode from blocking - // until ffmpeg has emitted its first bytes. - stream io.Reader -} - -func (f *flacPatcher) Read(p []byte) (int, error) { - if f.stream == nil { - prefix := make([]byte, flacPrefixLen) - n, err := io.ReadFull(f.ReadCloser, prefix) - if err != nil && !errors.Is(err, io.EOF) && !errors.Is(err, io.ErrUnexpectedEOF) { - return 0, err - } - prefix = prefix[:n] - if err == nil { - setFLACTotalSamples(prefix, f.duration) - } - f.stream = io.MultiReader(bytes.NewReader(prefix), f.ReadCloser) - } - return f.stream.Read(p) + setFLACTotalSamples(buf, h.duration) + return buf, nil } // setFLACTotalSamples takes the rate from the header rather than the transcode // options, so a resampled (-ar) output still gets the right count. func setFLACTotalSamples(prefix []byte, duration float32) { - if string(prefix[:4]) != "fLaC" || prefix[4]&0x7F != 0 { + if prefix[4]&0x7F != 0 { // the first metadata block must be STREAMINFO return } // 20-bit rate | 3-bit channels | 5-bit depth | 36-bit total_samples diff --git a/core/ffmpeg/flac_streaminfo_test.go b/core/ffmpeg/flac_streaminfo_test.go index 6bf3503d7..741f74b53 100644 --- a/core/ffmpeg/flac_streaminfo_test.go +++ b/core/ffmpeg/flac_streaminfo_test.go @@ -31,7 +31,7 @@ var _ = Describe("patchFLACDuration", func() { } readAll := func(in []byte, duration float32) []byte { - out, err := io.ReadAll(patchFLACDuration(io.NopCloser(bytes.NewReader(in)), duration)) + out, err := io.ReadAll(patchPipedHeader(io.NopCloser(bytes.NewReader(in)), duration)) Expect(err).ToNot(HaveOccurred()) return out } @@ -118,14 +118,14 @@ var _ = Describe("patchFLACDuration", func() { }) It("propagates a read error from the underlying stream", func() { - _, err := io.ReadAll(patchFLACDuration(io.NopCloser(io.MultiReader( + _, err := io.ReadAll(patchPipedHeader(io.NopCloser(io.MultiReader( bytes.NewReader(pipedFLAC()[:10]), &errReader{})), 1.0)) Expect(err).To(MatchError("boom")) }) It("closes the underlying stream", func() { c := &closeSpy{Reader: bytes.NewReader(pipedFLAC())} - Expect(patchFLACDuration(c, 1.0).Close()).To(Succeed()) + Expect(patchPipedHeader(c, 1.0).Close()).To(Succeed()) Expect(c.closed).To(BeTrue()) }) }) diff --git a/core/ffmpeg/mp3_xing.go b/core/ffmpeg/mp3_xing.go index 125469605..fd2c8524b 100644 --- a/core/ffmpeg/mp3_xing.go +++ b/core/ffmpeg/mp3_xing.go @@ -1,52 +1,24 @@ package ffmpeg import ( - "bytes" "encoding/binary" - "errors" - "io" "math" "slices" + + "github.com/navidrome/navidrome/utils/gg" ) const ( mp3HeaderLen = 4 mp3ID3Len = 10 - mp3MaxPrefix = 1 << 20 // an ID3 tag carrying cover art still fits + mp3MaxPrefix = 64 << 10 // ffmpeg cannot write an attached picture to a pipe, so the tag stays small ) -// patchMP3Duration prepends the Xing frame that ffmpeg omits when writing to a pipe, -// since it only knows the frame count once it can rewind to the first frame. -func patchMP3Duration(r io.ReadCloser, duration float32) io.ReadCloser { - if duration <= 0 { - return r - } - return &mp3Patcher{ReadCloser: r, duration: duration} -} - -type mp3Patcher struct { - io.ReadCloser - duration float32 - // Peeking here rather than in the constructor keeps Transcode from blocking - // until ffmpeg has emitted its first bytes. - stream io.Reader -} - -func (m *mp3Patcher) Read(p []byte) (int, error) { - if m.stream == nil { - prefix, err := m.peek() - if err != nil { - return 0, err - } - m.stream = io.MultiReader(bytes.NewReader(prefix), m.ReadCloser) - } - return m.stream.Read(p) -} - -// peek returns the head of the stream with a Xing frame inserted before the first -// audio frame, or unchanged when it cannot make sense of what ffmpeg wrote. -func (m *mp3Patcher) peek() ([]byte, error) { - buf, err := m.fill(nil, mp3ID3Len) +// mp3Prefix returns the head of the stream with a Xing frame inserted before the first +// audio frame, which ffmpeg omits on a pipe since it only knows the frame count once it +// can rewind. It returns the head unchanged when it cannot make sense of what ffmpeg wrote. +func (h *headerPatcher) mp3Prefix(buf []byte) ([]byte, error) { + buf, err := h.fill(buf, mp3ID3Len) if err != nil || len(buf) < mp3ID3Len { return buf, err } @@ -56,39 +28,26 @@ func (m *mp3Patcher) peek() ([]byte, error) { return buf, nil } } - if buf, err = m.fill(buf, start+mp3HeaderLen); err != nil || len(buf) < start+mp3HeaderLen { + if buf, err = h.fill(buf, start+mp3HeaderLen); err != nil || len(buf) < start+mp3HeaderLen { return buf, err } frame, ok := parseMP3Header(buf[start:]) if !ok { return buf, nil } - if buf, err = m.fill(buf, start+frame.size); err != nil || len(buf) < start+frame.size { + if buf, err = h.fill(buf, start+frame.size); err != nil || len(buf) < start+frame.size { return buf, err } if isXingFrame(buf[start:], frame.tagOffset) { return buf, nil } - xing, ok := xingFrame(buf[start:], frame, m.duration) + xing, ok := xingFrame(buf[start:], frame, h.duration) if !ok { return buf, nil } return slices.Insert(buf, start, xing...), nil } -// fill grows buf to n bytes, stopping short when the stream ends first. -func (m *mp3Patcher) fill(buf []byte, n int) ([]byte, error) { - if len(buf) >= n { - return buf, nil - } - more := make([]byte, n-len(buf)) - read, err := io.ReadFull(m.ReadCloser, more) - if err != nil && !errors.Is(err, io.EOF) && !errors.Is(err, io.ErrUnexpectedEOF) { - return buf, err - } - return append(buf, more[:read]...), nil -} - func id3TagLen(header []byte) int { size := int(header[6]&0x7F)<<21 | int(header[7]&0x7F)<<14 | int(header[8]&0x7F)<<7 | int(header[9]&0x7F) if header[5]&0x10 != 0 { // footer present @@ -102,7 +61,7 @@ type mp3Frame struct { samples int // per frame size int // bytes, including the header sideInfo int - tagOffset int // where a Xing tag would sit in this frame + tagOffset int // where a Xing tag sits in the incoming frame, which may carry a CRC } var ( @@ -125,7 +84,7 @@ func parseMP3Header(h []byte) (mp3Frame, bool) { } mpeg1 := version == 3 sampleRate := mp3SampleRates[version][(h[2]>>2)&0x03] - bitRate := 1000 * mp3BitRates[b2i(!mpeg1)][h[2]>>4] + bitRate := 1000 * mp3BitRates[gg.If(mpeg1, 0, 1)][h[2]>>4] if sampleRate == 0 || bitRate == 0 { return mp3Frame{}, false } @@ -173,10 +132,3 @@ func xingFrame(first []byte, f mp3Frame, duration float32) ([]byte, bool) { binary.BigEndian.PutUint32(frame[at+8:], uint32(frames)) return frame, true } - -func b2i(b bool) int { - if b { - return 1 - } - return 0 -} diff --git a/core/ffmpeg/mp3_xing_test.go b/core/ffmpeg/mp3_xing_test.go index 7d66268c2..2a66454c8 100644 --- a/core/ffmpeg/mp3_xing_test.go +++ b/core/ffmpeg/mp3_xing_test.go @@ -34,26 +34,25 @@ var _ = Describe("patchMP3Duration", func() { }) readAll := func(in []byte, duration float32) []byte { - out, err := io.ReadAll(patchMP3Duration(io.NopCloser(bytes.NewReader(in)), duration)) + out, err := io.ReadAll(patchPipedHeader(io.NopCloser(bytes.NewReader(in)), duration)) Expect(err).ToNot(HaveOccurred()) return out } - // Reads the tag, flags and frame count of the frame starting at frameStart. - readXing := func(b []byte, frameStart, tagOffset int) (string, uint32, uint32) { + // Reads the tag and frame count of the frame starting at frameStart. + readXing := func(b []byte, frameStart, tagOffset int) (string, uint32) { at := frameStart + tagOffset - return string(b[at : at+4]), - binary.BigEndian.Uint32(b[at+4:]), - binary.BigEndian.Uint32(b[at+8:]) + return string(b[at : at+4]), binary.BigEndian.Uint32(b[at+8:]) } It("inserts a Xing frame declaring the duration in frames", func() { out := readAll(pipedMP3, 1.0) - tag, flags, frames := readXing(out, fixtureID3Len, stereoTagOffset) + tag, frames := readXing(out, fixtureID3Len, stereoTagOffset) Expect(tag).To(Equal("Xing")) - Expect(flags).To(Equal(uint32(1)), "only the frame count is present") Expect(frames).To(Equal(uint32(38)), "1s at 44100Hz is 38 frames of 1152 samples") + flags := binary.BigEndian.Uint32(out[fixtureID3Len+stereoTagOffset+4:]) + Expect(flags).To(Equal(uint32(1)), "only the frame count is present") }) It("leaves the ID3 tag and the audio frames untouched", func() { @@ -72,13 +71,13 @@ var _ = Describe("patchMP3Duration", func() { out := readAll(in, 1.0) - _, _, frames := readXing(out, start, stereoTagOffset) + _, frames := readXing(out, start, stereoTagOffset) Expect(frames).To(Equal(uint32(42))) // 48000/1152, rounded }) It("rounds the frame count to the nearest frame", func() { out := readAll(pipedMP3, 0.99) // 37.9 frames - _, _, frames := readXing(out, fixtureID3Len, stereoTagOffset) + _, frames := readXing(out, fixtureID3Len, stereoTagOffset) Expect(frames).To(Equal(uint32(38))) }) @@ -87,7 +86,7 @@ var _ = Describe("patchMP3Duration", func() { out := readAll(in, 1.0) - tag, _, frames := readXing(out, 0, monoTagOffset) + tag, frames := readXing(out, 0, monoTagOffset) Expect(tag).To(Equal("Xing")) Expect(frames).To(Equal(uint32(38))) }) @@ -97,7 +96,7 @@ var _ = Describe("patchMP3Duration", func() { out := readAll(in, 1.0) - tag, _, _ := readXing(out, 0, stereoTagOffset) + tag, _ := readXing(out, 0, stereoTagOffset) Expect(tag).To(Equal("Xing")) Expect(out[fixtureFrameLen:]).To(Equal(in)) }) @@ -142,16 +141,15 @@ var _ = Describe("patchMP3Duration", func() { }) It("propagates a read error from the underlying stream", func() { - r := patchMP3Duration(io.NopCloser(iotest.ErrReader(errors.New("ffmpeg died"))), 1.0) + r := patchPipedHeader(io.NopCloser(iotest.ErrReader(errors.New("ffmpeg died"))), 1.0) _, err := io.ReadAll(r) Expect(err).To(MatchError(ContainSubstring("ffmpeg died"))) }) }) -// monoMP3 is two silent MPEG1 Layer III frames, 128kbps, 44100Hz, mono. +// monoMP3 is one silent MPEG1 Layer III frame, 128kbps, 44100Hz, mono. func monoMP3() []byte { - b := make([]byte, 2*417) + b := make([]byte, 417) copy(b, []byte{0xFF, 0xFB, 0x90, 0xC0}) - copy(b[417:], []byte{0xFF, 0xFB, 0x90, 0xC0}) return b } diff --git a/core/ffmpeg/piped_header.go b/core/ffmpeg/piped_header.go new file mode 100644 index 000000000..71f9ee754 --- /dev/null +++ b/core/ffmpeg/piped_header.go @@ -0,0 +1,61 @@ +package ffmpeg + +import ( + "bytes" + "errors" + "io" +) + +// patchPipedHeader repairs the header ffmpeg cannot finish when its output is a pipe, +// since it only learns the sample or frame count once it can no longer rewind. +func patchPipedHeader(r io.ReadCloser, duration float32) io.ReadCloser { + if duration <= 0 { + return r + } + return &headerPatcher{ReadCloser: r, duration: duration} +} + +type headerPatcher struct { + io.ReadCloser + duration float32 + // Peeking here rather than in the constructor keeps Transcode from blocking + // until ffmpeg has emitted its first bytes. + stream io.Reader +} + +func (h *headerPatcher) Read(p []byte) (int, error) { + if h.stream == nil { + prefix, err := h.peek() + if err != nil { + return 0, err + } + h.stream = io.MultiReader(bytes.NewReader(prefix), h.ReadCloser) + } + return h.stream.Read(p) +} + +// peek dispatches on the magic bytes rather than on the requested format, which is a +// free-text label on the transcoding profile and need not match what ffmpeg emits. +func (h *headerPatcher) peek() ([]byte, error) { + buf, err := h.fill(nil, 4) + if err != nil || len(buf) < 4 { + return buf, err + } + if string(buf) == "fLaC" { + return h.flacPrefix(buf) + } + return h.mp3Prefix(buf) +} + +// fill grows buf to n bytes, stopping short when the stream ends first. +func (h *headerPatcher) fill(buf []byte, n int) ([]byte, error) { + if len(buf) >= n { + return buf, nil + } + more := make([]byte, n-len(buf)) + read, err := io.ReadFull(h.ReadCloser, more) + if err != nil && !errors.Is(err, io.EOF) && !errors.Is(err, io.ErrUnexpectedEOF) { + return buf, err + } + return append(buf, more[:read]...), nil +}