diff --git a/core/ffmpeg/ffmpeg.go b/core/ffmpeg/ffmpeg.go index af59178af..66760c461 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 header. + Duration float32 // seconds; 0 = unknown. Only used to repair a header ffmpeg piped out unfinished. } // AudioProbeResult contains authoritative audio stream properties from ffprobe. @@ -89,7 +89,7 @@ func (e *ffmpeg) Transcode(ctx context.Context, opts TranscodeOptions) (io.ReadC if err != nil { return nil, err } - return patchFLACDuration(out, opts.Duration-float32(opts.Offset)), 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/ffmpeg_test.go b/core/ffmpeg/ffmpeg_test.go index 46684fe14..bc517da78 100644 --- a/core/ffmpeg/ffmpeg_test.go +++ b/core/ffmpeg/ffmpeg_test.go @@ -1,7 +1,9 @@ package ffmpeg import ( + "bytes" "context" + "encoding/binary" "errors" "io" "os" @@ -712,6 +714,23 @@ var _ = Describe("ffmpeg", func() { Expect(err).ToNot(HaveOccurred()) Expect(readTotalSamples(out)).To(Equal(uint64(2 * 44100))) }) + + It("inserts an Info frame on a piped mp3 transcode", func() { + stream, err := ff.Transcode(GinkgoT().Context(), TranscodeOptions{ + Command: "ffmpeg -i %s -map 0:a:0 -v 0 -b:a 128k -f mp3 -", + Format: "mp3", + FilePath: "tests/fixtures/test.flac", + Duration: 1, // the fixture is exactly 1s at 44100Hz + }) + Expect(err).ToNot(HaveOccurred()) + defer stream.Close() + + out, err := io.ReadAll(stream) + Expect(err).ToNot(HaveOccurred()) + at := bytes.Index(out, []byte("Info")) + Expect(at).To(BeNumerically(">", 0)) + Expect(binary.BigEndian.Uint32(out[at+8:])).To(Equal(uint32(38))) + }) }) Context("stderr capture", func() { 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 new file mode 100644 index 000000000..a173cf671 --- /dev/null +++ b/core/ffmpeg/mp3_xing.go @@ -0,0 +1,143 @@ +package ffmpeg + +import ( + "encoding/binary" + "math" + "slices" + + "github.com/navidrome/navidrome/utils/gg" +) + +const ( + mp3HeaderLen = 4 + mp3ID3Len = 10 + mp3InfoTagLen = 12 // tag, flags and frame count + mp3MaxPrefix = 64 << 10 // ffmpeg cannot write an attached picture to a pipe, so the tag stays small +) + +// mp3Prefix returns the head of the stream with an Info 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 + } + start := 0 + if string(buf[:3]) == "ID3" { + if start = id3TagLen(buf); start > mp3MaxPrefix { + return buf, nil + } + } + 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 = 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:], h.duration) + if !ok { + return buf, nil + } + return slices.Insert(buf, start, xing...), 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 + size += mp3ID3Len + } + return mp3ID3Len + size +} + +type mp3Frame struct { + sampleRate int + samples int // per frame + size int // bytes, including the header + tagOffset int // where a Xing tag sits in the incoming frame, which may carry a CRC +} + +var ( + mp3BitRates = [2][16]int{ + {0, 32, 40, 48, 56, 64, 80, 96, 112, 128, 160, 192, 224, 256, 320, 0}, // MPEG 1 + {0, 8, 16, 24, 32, 40, 48, 56, 64, 80, 96, 112, 128, 144, 160, 0}, // MPEG 2 and 2.5 + } + mp3SampleRates = [4][4]int{ + {11025, 12000, 8000, 0}, // MPEG 2.5 + {}, // reserved + {22050, 24000, 16000, 0}, // MPEG 2 + {44100, 48000, 32000, 0}, // MPEG 1 + } +) + +func parseMP3Header(h []byte) (mp3Frame, bool) { + version, layer := (h[1]>>3)&0x03, (h[1]>>1)&0x03 + if h[0] != 0xFF || h[1]&0xE0 != 0xE0 || version == 1 || layer != 1 { + return mp3Frame{}, false // not the sync word of an MPEG Layer III frame + } + mpeg1 := version == 3 + sampleRate := mp3SampleRates[version][(h[2]>>2)&0x03] + bitRate := 1000 * mp3BitRates[gg.If(mpeg1, 0, 1)][h[2]>>4] + if sampleRate == 0 || bitRate == 0 { + return mp3Frame{}, false + } + f := mp3Frame{sampleRate: sampleRate, samples: 576} + var sideInfo int + mono := (h[3]>>6)&0x03 == 3 + switch { + case mpeg1 && mono: + f.samples, sideInfo = 1152, 17 + case mpeg1: + f.samples, sideInfo = 1152, 32 + case mono: + sideInfo = 9 + default: + sideInfo = 17 + } + f.size = f.samples/8*bitRate/sampleRate + int((h[2]>>1)&0x01) + f.tagOffset = mp3HeaderLen + sideInfo + if h[1]&0x01 == 0 { // CRC follows the header + f.tagOffset += 2 + } + return f, true +} + +func isXingFrame(frame []byte, tagOffset int) bool { + if len(frame) < tagOffset+4 { + return false + } + tag := string(frame[tagOffset : tagOffset+4]) + return tag == "Xing" || tag == "Info" +} + +// xingFrame builds a silent Info frame declaring how many frames follow it. Its header is +// the first frame's minus the CRC, with the bitrate raised until the frame fits the tag. +func xingFrame(first []byte, duration float32) ([]byte, bool) { + header := [mp3HeaderLen]byte(first) + header[1] |= 0x01 + f, ok := parseMP3Header(header[:]) + for ok && f.size < f.tagOffset+mp3InfoTagLen { + header[2] += 0x10 // next bitrate index; 15 is invalid, which ends the loop + f, ok = parseMP3Header(header[:]) + } + if !ok { + return nil, false + } + frames := math.Round(float64(duration) * float64(f.sampleRate) / float64(f.samples)) + if frames < 1 || frames > math.MaxUint32 { + return nil, false + } + frame := make([]byte, f.size) + copy(frame, header[:]) + copy(frame[f.tagOffset:], "Info") // a Xing tag without a seek table marks the stream unseekable to some decoders + binary.BigEndian.PutUint32(frame[f.tagOffset+4:], 1) // only the frame count is present + binary.BigEndian.PutUint32(frame[f.tagOffset+8:], uint32(frames)) + return frame, true +} diff --git a/core/ffmpeg/mp3_xing_test.go b/core/ffmpeg/mp3_xing_test.go new file mode 100644 index 000000000..d3bf8ff1a --- /dev/null +++ b/core/ffmpeg/mp3_xing_test.go @@ -0,0 +1,170 @@ +package ffmpeg + +import ( + "bytes" + "encoding/binary" + "errors" + "io" + "os" + "testing/iotest" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// Offsets of the Xing tag inside the first frame, for the two layouts the specs use. +const ( + stereoTagOffset = 36 // 4 header + 32 side info + monoTagOffset = 21 // 4 header + 17 side info + fixtureID3Len = 39581 // where the first frame of tests/fixtures/test.mp3 starts + fixtureFrameLen = 627 // its frame size, at 192kbps and 44100Hz +) + +var _ = Describe("patchMP3Duration", func() { + var pipedMP3 []byte + + // The fixture is an ffmpeg pipe's output in every way that matters here: ID3 tag, + // frames, no Xing. + BeforeEach(func() { + var err error + pipedMP3, err = os.ReadFile("tests/fixtures/test.mp3") + Expect(err).ToNot(HaveOccurred()) + Expect(pipedMP3[fixtureID3Len : fixtureID3Len+2]).To(Equal([]byte{0xFF, 0xFB})) + Expect(bytes.Contains(pipedMP3[:4096], []byte("Xing"))).To(BeFalse()) + Expect(bytes.Contains(pipedMP3[:4096], []byte("Info"))).To(BeFalse()) + }) + + readAll := func(in []byte, duration float32) []byte { + out, err := io.ReadAll(patchPipedHeader(io.NopCloser(bytes.NewReader(in)), duration)) + Expect(err).ToNot(HaveOccurred()) + return out + } + + // 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+8:]) + } + + It("inserts an Info frame declaring the duration in frames", func() { + out := readAll(pipedMP3, 1.0) + + tag, frames := readXing(out, fixtureID3Len, stereoTagOffset) + Expect(tag).To(Equal("Info")) + 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() { + out := readAll(pipedMP3, 1.0) + + start := fixtureID3Len + Expect(out).To(HaveLen(len(pipedMP3) + fixtureFrameLen)) + Expect(out[:start]).To(Equal(pipedMP3[:start])) + Expect(out[start+fixtureFrameLen:]).To(Equal(pipedMP3[start:])) + }) + + It("takes the sample rate from the header, not from the source file", func() { + in := bytes.Clone(pipedMP3) + start := fixtureID3Len + in[start+2] = in[start+2]&^0x0C | 0x04 // sample rate index 1: 48000Hz + + out := readAll(in, 1.0) + + _, 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) + Expect(frames).To(Equal(uint32(38))) + }) + + It("places the tag after the shorter side info of a mono stream", func() { + in := monoMP3() + + out := readAll(in, 1.0) + + tag, frames := readXing(out, 0, monoTagOffset) + Expect(tag).To(Equal("Info")) + Expect(frames).To(Equal(uint32(38))) + }) + + It("raises the bitrate of the inserted frame when the first frame is too small for the tag", func() { + // MPEG 2, 8kbps, 22050Hz, stereo: 26 bytes, as a VBR encode starting on silence emits. + in := make([]byte, 26) + copy(in, []byte{0xFF, 0xF3, 0x10, 0x00}) + + out := readAll(in, 1.0) + + Expect(out[2]>>4).To(Equal(byte(2)), "16kbps, the first bitrate whose frame fits the tag") + tag, frames := readXing(out, 0, 21) // 4 header + 17 side info + Expect(tag).To(Equal("Info")) + Expect(frames).To(Equal(uint32(38))) // 1s at 22050Hz is 38 frames of 576 samples + Expect(out[52:]).To(Equal(in)) // 52 bytes = frame size at 16kbps, 22050Hz + }) + + It("inserts the frame at the start of a stream with no ID3 tag", func() { + in := pipedMP3[fixtureID3Len:] + + out := readAll(in, 1.0) + + tag, _ := readXing(out, 0, stereoTagOffset) + Expect(tag).To(Equal("Info")) + Expect(out[fixtureFrameLen:]).To(Equal(in)) + }) + + It("leaves a stream that already declares its duration alone", func() { + patched := readAll(pipedMP3, 1.0) + Expect(readAll(patched, 99.0)).To(Equal(patched)) + }) + + It("passes through when the duration is zero or negative", func() { + Expect(readAll(pipedMP3, 0)).To(Equal(pipedMP3)) + Expect(readAll(pipedMP3, -5)).To(Equal(pipedMP3)) + }) + + It("passes through a stream that is not mp3", func() { + in := []byte("fLaC\x00\x00\x00\x22 and then some bytes that are not frames") + Expect(readAll(in, 1.0)).To(Equal(in)) + }) + + It("passes through a frame header with a reserved bitrate or sample rate", func() { + in := monoMP3() + in[2] |= 0xF0 // bitrate index 15 is invalid + Expect(readAll(in, 1.0)).To(Equal(in)) + + in = monoMP3() + in[2] |= 0x0C // sample rate index 3 is reserved + Expect(readAll(in, 1.0)).To(Equal(in)) + }) + + It("passes through a stream too short to hold a frame header", func() { + in := monoMP3()[:3] + Expect(readAll(in, 1.0)).To(Equal(in)) + }) + + It("passes through a stream whose ID3 tag never ends", func() { + in := bytes.Clone(pipedMP3[:200]) + Expect(readAll(in, 1.0)).To(Equal(in)) + }) + + It("passes through an empty stream", func() { + Expect(readAll(nil, 1.0)).To(BeEmpty()) + }) + + It("propagates a read error from the underlying stream", func() { + 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 one silent MPEG1 Layer III frame, 128kbps, 44100Hz, mono. +func monoMP3() []byte { + b := make([]byte, 417) + copy(b, []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 +} diff --git a/core/stream/media_streamer.go b/core/stream/media_streamer.go index c250b0c4e..14dd77b3f 100644 --- a/core/stream/media_streamer.go +++ b/core/stream/media_streamer.go @@ -146,6 +146,9 @@ func (s *Stream) Duration() float32 { return s.mf.Duration } func (s *Stream) ContentType() string { return mime.TypeByExtension("." + s.format) } func (s *Stream) Name() string { return s.mf.Title + "." + s.format } func (s *Stream) ModTime() time.Time { return s.mf.UpdatedAt } + +// EstimatedContentLength deliberately overshoots by 2.4%: the body also carries tags and +// padding, and net/http aborts a response that writes past its declared length. func (s *Stream) EstimatedContentLength() int { return int(s.mf.Duration * float32(s.bitRate) / 8 * 1024) } diff --git a/core/stream/media_streamer_test.go b/core/stream/media_streamer_test.go index 5a4bcd480..c3652379b 100644 --- a/core/stream/media_streamer_test.go +++ b/core/stream/media_streamer_test.go @@ -8,6 +8,7 @@ import ( "net/http" "net/http/httptest" "os" + "strconv" "testing/iotest" "time" @@ -183,6 +184,20 @@ var _ = Describe("MediaStreamer", func() { _, err = io.ReadAll(resp.Body) Expect(err).To(HaveOccurred()) }) + + It("estimates a content length above the nominal size, so the body never outruns it", func() { + hundredSeconds := *mf + hundredSeconds.Duration = 100 + s := stream.NewStream(&hundredSeconds, "mp3", 128, io.NopCloser(bytes.NewReader(nil))) + w := httptest.NewRecorder() + r := httptest.NewRequest(http.MethodGet, "/?estimateContentLength=true", nil) + + _, _ = s.Serve(ctx, w, r) + + length, err := strconv.Atoi(w.Header().Get("Content-Length")) + Expect(err).ToNot(HaveOccurred()) + Expect(length).To(BeNumerically(">", 100*128*1000/8)) + }) }) })