diff --git a/go/audio/amplification.go b/go/audio/amplification.go index 66cd235..f973428 100644 --- a/go/audio/amplification.go +++ b/go/audio/amplification.go @@ -15,50 +15,31 @@ import ( // (pcm_f32le, audioFormat=3). The write is atomic: samples are written to a // temp file in the same directory and then renamed over the original. func Gain(path string, offsetDB float64) error { - f, err := os.Open(path) + left, right, sampleRate, err := readSamples(path) if err != nil { - return fmt.Errorf("opening %s: %w", path, err) - } - - h, err := readWAVHeader(f) - if err != nil { - f.Close() - return fmt.Errorf("reading WAV header: %w", err) - } - - raw := make([]byte, h.dataSize) - if _, err := io.ReadFull(f, raw); err != nil { - f.Close() - return fmt.Errorf("reading PCM data: %w", err) - } - f.Close() - - var left, right []float64 - switch { - case h.audioFormat == 3 && h.bitsPerSample == 32: - left, right = SamplesFromFloat32(raw) - case h.audioFormat == 1 && h.bitsPerSample == 24: - left, right = SamplesFromInt24(raw) - case h.audioFormat == 1 && h.bitsPerSample == 32: - left, right = SamplesFromInt32(raw) - case h.audioFormat == 1 && h.bitsPerSample == 16: - left, right = SamplesFromInt16(raw) - default: - return fmt.Errorf("unsupported format: audioformat=%d bits=%d", h.audioFormat, h.bitsPerSample) + return fmt.Errorf("gain: %w", err) } + GainSamples(left, right, offsetDB) + return WriteFloat32WAV(path, sampleRate, left, right) +} +// GainSamples applies a linear gain of offsetDB decibels to both channels in +// place. It is the in-memory core of Gain. +func GainSamples(left, right []float64, offsetDB float64) { gain := math.Pow(10, offsetDB/20) for i := range left { left[i] *= gain + } + for i := range right { right[i] *= gain } - - return writeFloat32WAV(path, h.sampleRate, left, right) } -// writeFloat32WAV writes stereo float64 samples to a 32-bit float WAV at path, -// atomically via a temp file in the same directory followed by a rename. -func writeFloat32WAV(path string, sampleRate uint32, left, right []float64) error { +// WriteFloat32WAV writes stereo float64 samples to a 32-bit float WAV at path, +// atomically via a temp file in the same directory followed by a fsync and +// rename — the sync matters on servers, where a crash mid-pipeline must not +// leave a truncated file behind a completed rename. +func WriteFloat32WAV(path string, sampleRate uint32, left, right []float64) error { dir := filepath.Dir(path) tmp, err := os.CreateTemp(dir, ".gain-*.wav.tmp") if err != nil { @@ -71,6 +52,11 @@ func writeFloat32WAV(path string, sampleRate uint32, left, right []float64) erro os.Remove(tmpName) return err } + if err := tmp.Sync(); err != nil { + tmp.Close() + os.Remove(tmpName) + return fmt.Errorf("syncing temp file: %w", err) + } if err := tmp.Close(); err != nil { os.Remove(tmpName) return fmt.Errorf("closing temp file: %w", err) @@ -93,7 +79,15 @@ func writeFloat32WAVTo(w io.Writer, sampleRate uint32, left, right []float64) er blockAlign := uint16(numChannels * bitsPerSample / 8) byteRate := sampleRate * uint32(blockAlign) - dataSize := uint32(nFrames) * uint32(blockAlign) + + // RIFF sizes are uint32: past 4 GiB the header silently wraps and every + // downstream reader sees a corrupt file. Refuse instead — callers with + // longer material must split or lower the intermediate sample rate. + dataSize64 := uint64(nFrames) * uint64(blockAlign) + if dataSize64+36 > math.MaxUint32 { + return fmt.Errorf("WAV data would be %d bytes — over the 4 GiB RIFF limit", dataSize64) + } + dataSize := uint32(dataSize64) riffSize := 36 + dataSize // Buffer the writer: every header field and sample is emitted through bw as a diff --git a/go/audio/dynamics.go b/go/audio/dynamics.go index b76b00c..3811b89 100644 --- a/go/audio/dynamics.go +++ b/go/audio/dynamics.go @@ -2,10 +2,10 @@ package audio import ( "fmt" - "io" "math" - "os" + "runtime" "sort" + "sync" ) // This file holds three purely time-domain dynamics processors: @@ -82,7 +82,7 @@ func LookaheadLimiter(path string, thresholdDb, lookaheadMs, releaseMs float64) return fmt.Errorf("lookahead limiter: %w", err) } LookaheadLimitSamples(left, right, float64(sampleRate), thresholdDb, lookaheadMs, releaseMs) - return writeFloat32WAV(path, sampleRate, left, right) + return WriteFloat32WAV(path, sampleRate, left, right) } // LookaheadLimitSamples is the in-place core of LookaheadLimiter. @@ -114,18 +114,15 @@ func LookaheadLimitSamples(left, right []float64, sampleRate, thresholdDb, looka // reconstructed waveform — the real peaks between samples — not just the // discrete sample values. That is what makes this a genuine true-peak limiter: // the output's true peak is guaranteed under the ceiling at the working rate, - // no oversampled pipeline required to fake it. - tpL := truePeakEnvelope(left[:n]) - tpR := truePeakEnvelope(right[:n]) + // no oversampled pipeline required to fake it. The per-channel envelopes are + // folded straight into required[] rather than materialized — on multi-hour + // files the two extra float64 arrays were gigabytes of transient memory — + // and the scan is chunk-parallel: each output depends only on the 12 + // preceding (read-only) input samples and the writes are disjoint. required := make([]float64, n) - for i := range n { - tp := math.Max(tpL[i], tpR[i]) - if tp > thresh { - required[i] = thresh / tp - } else { - required[i] = 1 - } - } + chunkedParallel(n, func(start, end int) { + requiredGainRange(left, right, required, thresh, start, end) + }) // winMin[i] = min(required[i .. i+look]); the gain therefore starts dropping // as soon as a peak enters the lookahead window, never after it has passed. @@ -152,33 +149,67 @@ func LookaheadLimitSamples(left, right []float64, sampleRate, thresholdDb, looka } } -// truePeakEnvelope returns, per sample, the inter-sample (true) peak magnitude at -// that position — the max over the 4 BS.1770 polyphase sub-samples of the -// reconstructed waveform. It is the per-sample form of channelTruePeakLinear (same -// tpFIR), and it drives the true-peak limiter so it tames the peaks BETWEEN -// samples, not only those landing on them. -func truePeakEnvelope(samples []float64) []float64 { - n := len(samples) - env := make([]float64, n) +// requiredGainRange fills required[start:end] with the gain that tames the +// stereo-linked true peak at each position to thresh (1.0 where no reduction +// is needed). The FIR reads input back to start-11; zero before index 0. +func requiredGainRange(left, right, required []float64, thresh float64, start, end int) { const taps = 12 - for i := range samples { + for i := start; i < end; i++ { var mx float64 for phase := range 4 { - var acc float64 + var accL, accR float64 for k := range taps { idx := i - k if idx < 0 { break } - acc += tpFIR[phase][k] * samples[idx] + accL += tpFIR[phase][k] * left[idx] + accR += tpFIR[phase][k] * right[idx] + } + if a := math.Abs(accL); a > mx { + mx = a } - if a := math.Abs(acc); a > mx { + if a := math.Abs(accR); a > mx { mx = a } } - env[i] = mx + if mx > thresh { + required[i] = thresh / mx + } else { + required[i] = 1 + } + } +} + +// chunkedParallel splits [0, n) into one contiguous range per worker and +// runs fn over them concurrently. fn must be safe for disjoint ranges (the +// envelope/peak scans here only read shared input and write disjoint output). +// Small inputs run inline. +func chunkedParallel(n int, fn func(start, end int)) { + workers := runtime.GOMAXPROCS(0) + if workers > 8 { + workers = 8 + } + const minChunk = 1 << 16 + if n < minChunk*2 || workers < 2 { + fn(0, n) + return } - return env + chunk := (n + workers - 1) / workers + var wg sync.WaitGroup + for w := range workers { + start := w * chunk + if start >= n { + break + } + end := min(start+chunk, n) + wg.Add(1) + go func(start, end int) { + defer wg.Done() + fn(start, end) + }(start, end) + } + wg.Wait() } // slidingMin returns, for every index i, the minimum of x over the forward @@ -225,6 +256,13 @@ func MeasureDynamicsScore(path string) (*DynamicsScoreAnalysis, error) { return dynamicsScore(left, float64(sampleRate)), nil } +// MeasureDynamicsScoreSamples computes the Dynamics Score from in-memory +// channel-1 (left) samples — the buffer-level form of MeasureDynamicsScore for +// callers that keep the whole file resident across processing stages. +func MeasureDynamicsScoreSamples(left []float64, sampleRate float64) *DynamicsScoreAnalysis { + return dynamicsScore(left, sampleRate) +} + // dynamicsScore is the in-memory core of MeasureDynamicsScore, operating on the // channel-1 (left) samples — the channel astats' DS is parsed from. func dynamicsScore(x []float64, sampleRate float64) *DynamicsScoreAnalysis { @@ -336,7 +374,7 @@ func LimitConforming(path string, strength int) error { return fmt.Errorf("conformity limiter: %w", err) } LookaheadLimitSamples(left, right, float64(sampleRate), c.thresholdDb, c.attackMs, c.releaseMs) - return writeFloat32WAV(path, sampleRate, left, right) + return WriteFloat32WAV(path, sampleRate, left, right) } // LimitCharacter applies the soft-knee character limiter at the calibrated @@ -348,7 +386,7 @@ func LimitCharacter(path string, strength int) error { return fmt.Errorf("character limiter: %w", err) } CharacterLimitSamples(left, right, float64(sampleRate), c.thresholdDb, characterLimiterKnee, c.attackMs, c.releaseMs) - return writeFloat32WAV(path, sampleRate, left, right) + return WriteFloat32WAV(path, sampleRate, left, right) } // CharacterLimiter applies a soft-knee limiter at path in place, rewriting it as @@ -363,7 +401,7 @@ func CharacterLimiter(path string, thresholdDb, kneeDb, attackMs, releaseMs floa return fmt.Errorf("character limiter: %w", err) } CharacterLimitSamples(left, right, float64(sampleRate), thresholdDb, kneeDb, attackMs, releaseMs) - return writeFloat32WAV(path, sampleRate, left, right) + return WriteFloat32WAV(path, sampleRate, left, right) } // CharacterLimitSamples is the in-place core of CharacterLimiter: a feed-forward @@ -383,7 +421,7 @@ func Compress(path string, thresholdDb, ratio, kneeDb, attackMs, releaseMs, make return fmt.Errorf("compressor: %w", err) } CompressSamples(left, right, float64(sampleRate), thresholdDb, ratio, kneeDb, attackMs, releaseMs, makeupDb, rmsMs) - return writeFloat32WAV(path, sampleRate, left, right) + return WriteFloat32WAV(path, sampleRate, left, right) } // CompressSamples is the in-place core of Compress. @@ -432,7 +470,7 @@ func UpwardCompress(path string, thresholdDb, ratio, kneeDb, attackMs, releaseMs return fmt.Errorf("upward compressor: %w", err) } UpwardCompressSamples(left, right, float64(sampleRate), thresholdDb, ratio, kneeDb, attackMs, releaseMs, maxBoostDb, rmsMs) - return writeFloat32WAV(path, sampleRate, left, right) + return WriteFloat32WAV(path, sampleRate, left, right) } // UpwardCompressSamples is the in-place core of UpwardCompress. @@ -519,37 +557,9 @@ func processDynamics(left, right []float64, sampleRate, thresholdDb, ratio, knee } // readStereo reads a WAV at path into stereo float64 buffers and also returns -// its sample rate, decoding the same formats ReadWAV supports. ReadWAV itself -// does not surface the sample rate, which the time-based processors here need to -// turn millisecond time factors into per-sample coefficients. +// its sample rate. It is a thin alias for the package-central readSamples, +// kept for the time-based processors here that turn millisecond time factors +// into per-sample coefficients. func readStereo(path string) (left, right []float64, sampleRate uint32, err error) { - f, err := os.Open(path) - if err != nil { - return nil, nil, 0, err - } - defer f.Close() - - h, err := readWAVHeader(f) - if err != nil { - return nil, nil, 0, err - } - - raw := make([]byte, h.dataSize) - if _, err := io.ReadFull(f, raw); err != nil { - return nil, nil, 0, err - } - - switch { - case h.audioFormat == 3 && h.bitsPerSample == 32: - left, right = SamplesFromFloat32(raw) - case h.audioFormat == 1 && h.bitsPerSample == 24: - left, right = SamplesFromInt24(raw) - case h.audioFormat == 1 && h.bitsPerSample == 32: - left, right = SamplesFromInt32(raw) - case h.audioFormat == 1 && h.bitsPerSample == 16: - left, right = SamplesFromInt16(raw) - default: - return nil, nil, 0, fmt.Errorf("format not supported: audioformat: %d, bits per sample: %d", h.audioFormat, h.bitsPerSample) - } - return left, right, h.sampleRate, nil + return readSamples(path) } diff --git a/go/audio/dynaudnorm.go b/go/audio/dynaudnorm.go index d3483a5..ea61283 100644 --- a/go/audio/dynaudnorm.go +++ b/go/audio/dynaudnorm.go @@ -46,8 +46,10 @@ func BuildDynaudnormFilter(params *DynaudnormParams, isSpeech bool) string { return "" } + // gausssize must be odd — ffmpeg coerced the old 36 to 37 with a warning, + // so 37 is spelled out to keep the effective behaviour and lose the warning. return fmt.Sprintf( - "dynaudnorm=framelen=650:gausssize=36:targetrms=%.6f:threshold=%.6f:altboundary=true:overlap=0.95", + "dynaudnorm=framelen=650:gausssize=37:targetrms=%.6f:threshold=%.6f:altboundary=true:overlap=0.95", params.TargetRMS, params.Threshold, ) diff --git a/go/audio/lra.go b/go/audio/lra.go index fdc51af..63f63e1 100644 --- a/go/audio/lra.go +++ b/go/audio/lra.go @@ -10,44 +10,48 @@ package audio // This is the shared implementation used by both the app pipeline and the // audition CLI, so they can never drift. func ReduceLRA(path string, targetLRA float64, maxPasses int) (float64, error) { - r, err := MeasureLUFS(path) + left, right, sampleRate, err := readSamples(path) if err != nil { return 0, err } - lra := r.LRA + lra := ReduceLRASamples(left, right, float64(sampleRate), targetLRA, maxPasses) + if err := WriteFloat32WAV(path, sampleRate, left, right); err != nil { + return lra, err + } + return lra, nil +} + +// ReduceLRASamples is the in-memory core of ReduceLRA: it runs the same +// converging multipass directly on resident stereo buffers, so multi-pass +// processing costs one read and one write instead of ~ten file round-trips +// per pass. Returns the final LRA. +// +// The loop measures with loudnessRange directly — LRA is all it steers on, +// and skipping the integrated-loudness and true-peak legs of a full +// measurement saves most of each iteration's metering cost on long files. +func ReduceLRASamples(left, right []float64, sampleRate, targetLRA float64, maxPasses int) float64 { + lra := loudnessRange(left, right, sampleRate) for pass := 0; pass < maxPasses && lra > targetLRA; pass++ { before := lra - if err := reduceLRAPass(path, targetLRA); err != nil { - return lra, err - } - if r, err = MeasureLUFS(path); err != nil { - return lra, err - } - lra = r.LRA + reduceLRAPass(left, right, sampleRate, lra-targetLRA) + lra = loudnessRange(left, right, sampleRate) if before-lra < 0.3 { break // diminishing returns } } - return lra, nil + return lra } -// reduceLRAPass applies one gentle pass of the LRA chain, sized to the current -// overshoot. Thresholds come from the RMS statistics; all time constants are -// DS-driven via GetCompressionModifiers. -func reduceLRAPass(path string, targetLRA float64) error { - m, err := MeasureLUFS(path) - if err != nil { - return err - } - over := m.LRA - targetLRA +// reduceLRAPass applies one gentle pass of the LRA chain, sized to the +// current overshoot (over = current LRA − target, measured by the caller so +// it isn't re-measured here). Thresholds come from the RMS statistics; all +// time constants are DS-driven via GetCompressionModifiers. +func reduceLRAPass(left, right []float64, sampleRate, over float64) { if over <= 0 { - return nil + return } - ds, err := MeasureDynamicsScore(path) - if err != nil { - return err - } + ds := MeasureDynamicsScoreSamples(left, sampleRate) const wideKnee = 12.0 const rmsWindowMs = 300.0 // RMS detection — track sustained loudness, not peaks @@ -57,29 +61,22 @@ func reduceLRAPass(path string, targetLRA float64) error { interp := clamp(0.6-over*0.05, 0.1, 0.6) thresh := ds.RMSLevel + interp*(ds.RMSPeak-ds.RMSLevel) ratio := clamp(1+over*0.15, 1, 6) - if err := Compress(path, thresh, ratio, wideKnee, - 100*mods.AttackMultiplier, 1500*mods.ReleaseMultiplier, 0, rmsWindowMs); err != nil { - return err - } + CompressSamples(left, right, sampleRate, thresh, ratio, wideKnee, + 100*mods.AttackMultiplier, 1500*mods.ReleaseMultiplier, 0, rmsWindowMs) // Upward comp — lift passages below RMS level toward it. upRatio := clamp(1+over*0.1, 1, 4) const upMaxBoost = 5.0 - if err := UpwardCompress(path, ds.RMSLevel, upRatio, wideKnee, - 300*mods.AttackMultiplier, 600*mods.ReleaseMultiplier, upMaxBoost, rmsWindowMs); err != nil { - return err - } + UpwardCompressSamples(left, right, sampleRate, ds.RMSLevel, upRatio, wideKnee, + 300*mods.AttackMultiplier, 600*mods.ReleaseMultiplier, upMaxBoost, rmsWindowMs) // Slow character limiter — threshold from the post-comp DS, placed inside the // RMS-to-RMSpeak span and pushed down the more dynamic the material is. - ds2, err := MeasureDynamicsScore(path) - if err != nil { - return err - } + ds2 := MeasureDynamicsScoreSamples(left, sampleRate) span := ds2.RMSPeak - ds2.RMSLevel dyn := clamp((ds2.DynamicsScore-9)/12, 0, 1) charThresh := ds2.RMSPeak - (0.2+0.3*dyn)*span mods2 := GetCompressionModifiers(ds2.DynamicsScore) - return CharacterLimiter(path, charThresh, 6, + CharacterLimitSamples(left, right, sampleRate, charThresh, 6, 150*mods2.AttackMultiplier, 1200*mods2.ReleaseMultiplier) } diff --git a/go/audio/lufs.go b/go/audio/lufs.go index 9a11c1d..33029ef 100644 --- a/go/audio/lufs.go +++ b/go/audio/lufs.go @@ -2,10 +2,10 @@ package audio import ( "fmt" - "io" "math" - "os" + "runtime" "sort" + "sync" ) // biquadFilter is a single Direct-Form-I biquad IIR section. @@ -91,6 +91,22 @@ func kWeightChannel(samples []float64, sampleRate float64) []float64 { return out } +// kWeightStereo K-weights both channels concurrently. The biquads are +// IIR (serial per channel), so two goroutines — one per channel — is the +// available parallelism; on multi-hour files it halves the dominant cost of +// every loudness measurement. +func kWeightStereo(left, right []float64, sampleRate float64) (wL, wR []float64) { + var wg sync.WaitGroup + wg.Add(1) + go func() { + defer wg.Done() + wR = kWeightChannel(right, sampleRate) + }() + wL = kWeightChannel(left, sampleRate) + wg.Wait() + return wL, wR +} + // LUFSResult holds integrated loudness, true-peak and loudness-range, // per ITU-R BS.1770-5 and EBU Tech 3342. type LUFSResult struct { @@ -112,41 +128,22 @@ const ( // MeasureLUFS reads a WAV file and returns integrated loudness in LUFS. func MeasureLUFS(path string) (*LUFSResult, error) { - f, err := os.Open(path) - if err != nil { - return nil, err - } - defer f.Close() - - h, err := readWAVHeader(f) + left, right, sampleRate, err := readSamples(path) if err != nil { - return nil, fmt.Errorf("reading WAV header: %w", err) - } - - raw := make([]byte, h.dataSize) - if _, err := io.ReadFull(f, raw); err != nil { - return nil, fmt.Errorf("reading PCM data: %w", err) + return nil, fmt.Errorf("measure LUFS: %w", err) } + return MeasureLUFSSamples(left, right, float64(sampleRate)), nil +} - var left, right []float64 - switch { - case h.audioFormat == 3 && h.bitsPerSample == 32: - left, right = SamplesFromFloat32(raw) - case h.audioFormat == 1 && h.bitsPerSample == 24: - left, right = SamplesFromInt24(raw) - case h.audioFormat == 1 && h.bitsPerSample == 32: - left, right = SamplesFromInt32(raw) - case h.audioFormat == 1 && h.bitsPerSample == 16: - left, right = SamplesFromInt16(raw) - default: - return nil, fmt.Errorf("unsupported format: audioformat=%d bits=%d", h.audioFormat, h.bitsPerSample) +// MeasureLUFSSamples measures integrated loudness, true peak and LRA from +// in-memory stereo buffers — the buffer-level form of MeasureLUFS for callers +// that keep the whole file resident across processing stages. +func MeasureLUFSSamples(left, right []float64, sampleRate float64) *LUFSResult { + return &LUFSResult{ + Integrated: integratedLUFS(left, right, sampleRate), + TruePeak: truePeak(left, right), + LRA: loudnessRange(left, right, sampleRate), } - - sr := float64(h.sampleRate) - lufs := integratedLUFS(left, right, sr) - tp := truePeak(left, right) - lra := loudnessRange(left, right, sr) - return &LUFSResult{Integrated: lufs, TruePeak: tp, LRA: lra}, nil } // integratedLUFS implements the BS.1770-4 integrated loudness algorithm on @@ -157,8 +154,7 @@ func integratedLUFS(left, right []float64, sampleRate float64) float64 { } // Step 1: K-weight both channels. - wL := kWeightChannel(left, sampleRate) - wR := kWeightChannel(right, sampleRate) + wL, wR := kWeightStereo(left, right, sampleRate) // Step 2: Slice into 400 ms blocks with 100 ms hops and collect the // per-block power (sum of channel mean-squares, as BS.1770 specifies). @@ -271,9 +267,19 @@ var tpFIR = [4][12]float64{ // truePeak returns the inter-sample true-peak level in dBTP, taking the // highest across the left and right channels (ITU-R BS.1770-5 Annex 2). -// Returns math.Inf(-1) for a fully silent signal. +// Returns math.Inf(-1) for a fully silent signal. The two channels are +// scanned concurrently. func truePeak(left, right []float64) float64 { - pk := math.Max(channelTruePeakLinear(left), channelTruePeakLinear(right)) + var pkL, pkR float64 + var wg sync.WaitGroup + wg.Add(1) + go func() { + defer wg.Done() + pkR = channelTruePeakLinear(right) + }() + pkL = channelTruePeakLinear(left) + wg.Wait() + pk := math.Max(pkL, pkR) if pk <= 0 { return math.Inf(-1) } @@ -284,14 +290,57 @@ func truePeak(left, right []float64) float64 { // channelTruePeakLinear 4× over-samples one channel through the BS.1770-5 // polyphase FIR and returns the maximum absolute interpolated value (linear). +// +// The scan is chunk-parallel: each output sample depends only on the 12 +// preceding input samples, and the input is read-only, so workers can take +// disjoint index ranges with no coordination beyond the final max-reduce. +// This is the hottest loop in the package (96 multiply-adds per frame) — +// on a long file the speedup is essentially linear in cores. func channelTruePeakLinear(samples []float64) float64 { n := len(samples) if n == 0 { return 0 } + workers := runtime.GOMAXPROCS(0) + if workers > 8 { + workers = 8 + } + const minChunk = 1 << 16 + if n < minChunk*2 || workers < 2 { + return truePeakRange(samples, 0, n) + } + + chunk := (n + workers - 1) / workers + maxes := make([]float64, workers) + var wg sync.WaitGroup + for w := range workers { + start := w * chunk + if start >= n { + break + } + end := min(start+chunk, n) + wg.Add(1) + go func(w, start, end int) { + defer wg.Done() + maxes[w] = truePeakRange(samples, start, end) + }(w, start, end) + } + wg.Wait() + var maxAbs float64 + for _, m := range maxes { + if m > maxAbs { + maxAbs = m + } + } + return maxAbs +} + +// truePeakRange runs the polyphase FIR over output indices [start, end), +// reading input samples back to start-11 (zero before index 0). +func truePeakRange(samples []float64, start, end int) float64 { const taps = 12 var maxAbs float64 - for i := range samples { + for i := start; i < end; i++ { for phase := range 4 { var acc float64 for k := range taps { @@ -318,8 +367,7 @@ func loudnessRange(left, right []float64, sampleRate float64) float64 { return 0 } - wL := kWeightChannel(left, sampleRate) - wR := kWeightChannel(right, sampleRate) + wL, wR := kWeightStereo(left, right, sampleRate) blockSize := int(math.Round(lraBlockSizeSec * sampleRate)) hopSize := int(math.Round(lraHopSizeSec * sampleRate)) diff --git a/go/audio/read-pcm-bytes.go b/go/audio/read-pcm-bytes.go index 2077283..51936df 100644 --- a/go/audio/read-pcm-bytes.go +++ b/go/audio/read-pcm-bytes.go @@ -4,6 +4,7 @@ import ( "encoding/binary" "fmt" "io" + "math" "os" ) @@ -74,8 +75,8 @@ func readWAVHeader(r io.ReadSeeker) (wavHeader, error) { } } - if fmtChunk.NumChannels != 2 { - return wavHeader{}, fmt.Errorf("is not stereo, got %d channels", fmtChunk.NumChannels) + if fmtChunk.NumChannels != 1 && fmtChunk.NumChannels != 2 { + return wavHeader{}, fmt.Errorf("unsupported channel count: got %d, want mono or stereo", fmtChunk.NumChannels) } if format != 1 && format != 3 { return wavHeader{}, fmt.Errorf("audio format not supported: %d", format) @@ -90,6 +91,13 @@ func readWAVHeader(r io.ReadSeeker) (wavHeader, error) { if !fmtFound { return wavHeader{}, fmt.Errorf("data chunk before fmt chunk") } + // ffmpeg's wav muxer CLAMPS the size fields to uint32 max when + // the payload crosses the 4 GiB RIFF limit (and warns on + // stderr, but still exits 0). Trusting the clamped header would + // silently truncate hours of audio; refuse instead. + if chunkSize == 0xFFFFFFFF { + return wavHeader{}, fmt.Errorf("data chunk size is clamped at the 4 GiB RIFF limit — file is truncated or oversized") + } h.dataSize = chunkSize return h, nil @@ -108,23 +116,21 @@ func readWAVHeader(r io.ReadSeeker) (wavHeader, error) { } } -func ReadWAV(path string) (left, right []float64, err error) { - f, err := os.Open(path) - if err != nil { - return nil, nil, err - } - defer f.Close() - - h, err := readWAVHeader(f) - if err != nil { - return nil, nil, err - } - - raw := make([]byte, h.dataSize) - if _, err := io.ReadFull(f, raw); err != nil { - return nil, nil, err +// decodeSamples converts the raw data-chunk bytes described by h into +// per-channel float64 buffers. Stereo data decodes as-is; mono data is +// duplicated into two independent buffers so every downstream processor — +// all of which mutate left and right separately — behaves identically for +// mono and stereo input. This is the single decode switch for the whole +// package; ReadWAV, MeasureLUFS, Gain and readStereo all route through it. +func decodeSamples(h wavHeader, raw []byte) (left, right []float64, err error) { + if h.numChannels == 1 { + mono, err := decodeMono(h, raw) + if err != nil { + return nil, nil, err + } + right = append([]float64(nil), mono...) + return mono, right, nil } - switch { case h.audioFormat == 3 && h.bitsPerSample == 32: left, right = SamplesFromFloat32(raw) @@ -139,3 +145,88 @@ func ReadWAV(path string) (left, right []float64, err error) { } return left, right, nil } + +// decodeMono decodes single-channel PCM data into one float64 buffer, +// mirroring the format support of the stereo SamplesFrom* decoders. +func decodeMono(h wavHeader, raw []byte) ([]float64, error) { + switch { + case h.audioFormat == 3 && h.bitsPerSample == 32: + n := len(raw) / 4 + out := make([]float64, n) + for i := range n { + out[i] = float64(math.Float32frombits(binary.LittleEndian.Uint32(raw[i*4 : i*4+4]))) + } + return out, nil + case h.audioFormat == 1 && h.bitsPerSample == 24: + n := len(raw) / 3 + out := make([]float64, n) + const scale = 1.0 / 8388608.0 // 2^23 + for i := range n { + o := i * 3 + u := uint32(raw[o]) | uint32(raw[o+1])<<8 | uint32(raw[o+2])<<16 + s := int32(u) + if u&0x800000 != 0 { + s = int32(u | 0xFF000000) + } + out[i] = float64(s) * scale + } + return out, nil + case h.audioFormat == 1 && h.bitsPerSample == 32: + n := len(raw) / 4 + out := make([]float64, n) + const scale = 1.0 / 2147483648.0 // 2^31 + for i := range n { + out[i] = float64(int32(binary.LittleEndian.Uint32(raw[i*4:i*4+4]))) * scale + } + return out, nil + case h.audioFormat == 1 && h.bitsPerSample == 16: + n := len(raw) / 2 + out := make([]float64, n) + const scale = 1.0 / 32768.0 // 2^15 + for i := range n { + out[i] = float64(int16(binary.LittleEndian.Uint16(raw[i*2:i*2+2]))) * scale + } + return out, nil + default: + return nil, fmt.Errorf("format not supported: audioformat: %d, bits per sample: %d", h.audioFormat, h.bitsPerSample) + } +} + +// readSamples opens the WAV at path and returns decoded per-channel buffers +// plus the sample rate. Mono files come back as two duplicated buffers (see +// decodeSamples). +func readSamples(path string) (left, right []float64, sampleRate uint32, err error) { + f, err := os.Open(path) + if err != nil { + return nil, nil, 0, err + } + defer f.Close() + + h, err := readWAVHeader(f) + if err != nil { + return nil, nil, 0, err + } + + raw := make([]byte, h.dataSize) + if _, err := io.ReadFull(f, raw); err != nil { + return nil, nil, 0, err + } + left, right, err = decodeSamples(h, raw) + if err != nil { + return nil, nil, 0, err + } + return left, right, h.sampleRate, nil +} + +// ReadStereoWAV reads the WAV at path into stereo float64 buffers and returns +// its sample rate. Mono input is duplicated into both channels. This is the +// exported entry point for callers that drive the *Samples processors +// directly (servers, batch pipelines) instead of the path-based wrappers. +func ReadStereoWAV(path string) (left, right []float64, sampleRate uint32, err error) { + return readSamples(path) +} + +func ReadWAV(path string) (left, right []float64, err error) { + left, right, _, err = readSamples(path) + return left, right, err +} diff --git a/go/audio/wav_io_test.go b/go/audio/wav_io_test.go new file mode 100644 index 0000000..af235ad --- /dev/null +++ b/go/audio/wav_io_test.go @@ -0,0 +1,156 @@ +package audio + +import ( + "encoding/binary" + "math" + "os" + "path/filepath" + "testing" +) + +// writeTestWAV writes a minimal PCM WAV with the given channel count, bit +// depth (16 or 32-float) and interleaved samples in [-1, 1]. +func writeTestWAV(t *testing.T, path string, channels int, sampleRate uint32, interleaved []float64) { + t.Helper() + const bits = 16 + blockAlign := channels * bits / 8 + dataSize := len(interleaved) * bits / 8 + + buf := make([]byte, 0, 44+dataSize) + u16 := func(v uint16) { buf = binary.LittleEndian.AppendUint16(buf, v) } + u32 := func(v uint32) { buf = binary.LittleEndian.AppendUint32(buf, v) } + + buf = append(buf, "RIFF"...) + u32(uint32(36 + dataSize)) + buf = append(buf, "WAVE"...) + buf = append(buf, "fmt "...) + u32(16) + u16(1) // PCM int + u16(uint16(channels)) + u32(sampleRate) + u32(sampleRate * uint32(blockAlign)) + u16(uint16(blockAlign)) + u16(bits) + buf = append(buf, "data"...) + u32(uint32(dataSize)) + for _, s := range interleaved { + u16(uint16(int16(math.Round(s * 32767)))) + } + if err := os.WriteFile(path, buf, 0o644); err != nil { + t.Fatal(err) + } +} + +// TestReadStereoWAV_mono verifies a mono WAV decodes into two independent, +// identical channel buffers — so every stereo-linked processor treats mono +// input exactly like dual-mono stereo. +func TestReadStereoWAV_mono(t *testing.T) { + path := filepath.Join(t.TempDir(), "mono.wav") + mono := []float64{0, 0.25, -0.5, 1, -1, 0.125} + writeTestWAV(t, path, 1, 48000, mono) + + left, right, sr, err := ReadStereoWAV(path) + if err != nil { + t.Fatalf("ReadStereoWAV(mono): %v", err) + } + if sr != 48000 { + t.Fatalf("sample rate = %d, want 48000", sr) + } + if len(left) != len(mono) || len(right) != len(mono) { + t.Fatalf("lengths = %d/%d, want %d", len(left), len(right), len(mono)) + } + for i := range mono { + if math.Abs(left[i]-mono[i]) > 1.0/32000 { + t.Fatalf("left[%d] = %f, want ~%f", i, left[i], mono[i]) + } + if left[i] != right[i] { + t.Fatalf("channels differ at %d: %f vs %f", i, left[i], right[i]) + } + } + + // Buffers must be independent — in-place processors write both. + left[0] = 0.9 + if right[0] == 0.9 { + t.Fatal("left and right share backing memory") + } +} + +// TestReadStereoWAV_stereo guards the unchanged stereo path through the +// centralized decoder. +func TestReadStereoWAV_stereo(t *testing.T) { + path := filepath.Join(t.TempDir(), "stereo.wav") + interleaved := []float64{0.5, -0.5, 0.25, -0.25, 1, -1} + writeTestWAV(t, path, 2, 44100, interleaved) + + left, right, sr, err := ReadStereoWAV(path) + if err != nil { + t.Fatalf("ReadStereoWAV(stereo): %v", err) + } + if sr != 44100 { + t.Fatalf("sample rate = %d, want 44100", sr) + } + want := [][2]float64{{0.5, -0.5}, {0.25, -0.25}, {1, -1}} + if len(left) != len(want) { + t.Fatalf("frames = %d, want %d", len(left), len(want)) + } + for i, w := range want { + if math.Abs(left[i]-w[0]) > 1.0/32000 || math.Abs(right[i]-w[1]) > 1.0/32000 { + t.Fatalf("frame %d = %f/%f, want %f/%f", i, left[i], right[i], w[0], w[1]) + } + } +} + +// TestReduceLRASamples_matchesPath pins the samples-resident ReduceLRA to the +// path-based wrapper: same input, same result, same processed audio. +func TestReduceLRASamples_matchesPath(t *testing.T) { + const ( + sr = 8000 + secs = 30 + loudA = 0.5 + quietA = 0.05 + ) + // Alternate 3 s loud / 3 s quiet sine — wide LRA by construction. + n := sr * secs + left := make([]float64, n) + right := make([]float64, n) + for i := range n { + amp := loudA + if (i/(3*sr))%2 == 1 { + amp = quietA + } + // Quantize to float32 up front: the path-based run reads its input + // back from a float32 WAV, so the in-memory run must start from the + // same quantized data for the two to be bit-identical. + v := float64(float32(amp * math.Sin(2*math.Pi*440*float64(i)/sr))) + left[i] = v + right[i] = v + } + + pathLeft := append([]float64(nil), left...) + pathRight := append([]float64(nil), right...) + path := filepath.Join(t.TempDir(), "lra.wav") + if err := WriteFloat32WAV(path, sr, pathLeft, pathRight); err != nil { + t.Fatal(err) + } + + gotPath, err := ReduceLRA(path, 7, 2) + if err != nil { + t.Fatalf("ReduceLRA: %v", err) + } + gotSamples := ReduceLRASamples(left, right, sr, 7, 2) + + if math.Abs(gotPath-gotSamples) > 1e-9 { + t.Fatalf("path LRA %f != samples LRA %f", gotPath, gotSamples) + } + + outL, outR, _, err := ReadStereoWAV(path) + if err != nil { + t.Fatal(err) + } + for i := 0; i < n; i += 997 { // spot-check, float32 round-trip tolerance + if math.Abs(outL[i]-left[i]) > 1e-6 || math.Abs(outR[i]-right[i]) > 1e-6 { + t.Fatalf("processed audio diverges at %d: file %g/%g vs samples %g/%g", + i, outL[i], outR[i], left[i], right[i]) + } + } +} diff --git a/go/main.go b/go/main.go index 9199e05..5399243 100644 --- a/go/main.go +++ b/go/main.go @@ -1474,7 +1474,7 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool } } - if !noBitrateUsed { + if !noBitrateUsed && !n.noTranscode { if needsFullNumber { args = append(args, "-b:a", fmt.Sprintf("%d", bitrate)) } else { @@ -1489,7 +1489,7 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool args = append(args, "-application", "audio") } - usesDataCompression := actualCodec == "flac" || actualCodec == "libopus" + usesDataCompression := (actualCodec == "flac" || actualCodec == "libopus") && !n.noTranscode if usesDataCompression { var level int @@ -1594,8 +1594,11 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool cfg.EqTarget != "Off", !cfg.BypassProc)) - // Stage 1: EQ analysis (measures the raw input spectrum) - if cfg.EqTarget != "" && cfg.EqTarget != "Off" && !cfg.BypassProc { + // Stage 1: EQ analysis (measures the raw input spectrum). Every audio + // stage is skipped in no-transcode (tag-only) mode: the final render is + // -c copy, so a rendered intermediate could never reach the output — it + // would only be wasted work and a broken stream-copy source. + if cfg.EqTarget != "" && cfg.EqTarget != "Off" && !cfg.BypassProc && !n.noTranscode { eqBandAnalysis := n.analyzeFrequencyResponseBands(workingPath) if eqBandAnalysis == nil || len(eqBandAnalysis) == 0 { n.appLog.Write(fmt.Sprintf("✗ Failed to analyze frequency response: %s", filepath.Base(inputPath))) @@ -1651,7 +1654,7 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool !cfg.BypassProc)) var dsAnalysis *audio.DynamicsScoreAnalysis - if !cfg.BypassProc && (cfg.DynamicsPreset != "" && cfg.DynamicsPreset != "Off") { + if !cfg.BypassProc && !n.noTranscode && (cfg.DynamicsPreset != "" && cfg.DynamicsPreset != "Off") { dsAnalysis = n.calculateDynamicsScore(inputPath) if dsAnalysis == nil { n.appLog.Write(fmt.Sprintf("✗ Failed to calculate Dynamics Score: %s", filepath.Base(inputPath))) @@ -1660,7 +1663,7 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool } // Stage 2: Dynaudnorm (measures the post-EQ signal) - if cfg.DynNorm && !cfg.BypassProc { + if cfg.DynNorm && !cfg.BypassProc && !n.noTranscode { var dynaudnormFilter string dynamicsAnalysis := n.analyzeDynamics(workingPath) if dynamicsAnalysis == nil { @@ -1698,7 +1701,7 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool } // Stage 3: Compression (measures the post-EQ, post-dynaudnorm signal) - if cfg.DynamicsPreset != "" && cfg.DynamicsPreset != "Off" && !cfg.BypassProc { + if cfg.DynamicsPreset != "" && cfg.DynamicsPreset != "Off" && !cfg.BypassProc && !n.noTranscode { // MBC attenuates hot peaks before compressing var attenuatedPath string = workingPath @@ -1843,30 +1846,14 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool } } - type bitrateTiers struct { - mild, mid, hard int - } - - tiers := map[string]bitrateTiers{ - "libmp3lame": {mild: 192, mid: 160, hard: 128}, - "aac_at": {mild: 128, mid: 96, hard: 64}, - "libfdk_aac": {mild: 128, mid: 96, hard: 64}, - "libopus": {mild: 96, mid: 64, hard: 48}, - } - - // strength for brightness reduction - strength := 0 - - if t, ok := tiers[actualCodec]; ok { - switch { - case bitrate <= t.hard: - strength = 3 - case bitrate <= t.mid: - strength = 2 - case bitrate <= t.mild: - strength = 1 - } + // strength for brightness reduction and the encoder pre-limiter. The + // tiers are in kbps, but for the needsFullNumber codecs `bitrate` was + // converted to full bps above — convert back before comparing. + bitrateKbps := bitrate + if needsFullNumber { + bitrateKbps = bitrate / 1000 } + strength := brightnessStrength(actualCodec, bitrateKbps) n.logFile.Write("") n.logFile.Write(fmt.Sprintf("args: %s", args)) @@ -1886,44 +1873,55 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool cfg.IsSpeech, target, targetTp, ) - // LRA reduction toward the target range (self-gates if already under target). - const targetLRA = 7.0 - n.reduceLRA(workingPath, targetLRA, 4) + // Native loudness chain — rewrites workingPath IN PLACE as 32-bit float + // WAV. It runs even when cfg.BypassProc is set, on purpose: "Bypass all + // processing" is scoped to the EQ/Dynamics stages above (GUI: "Disables + // Dynamics and EQ regardless of selection above"); normalization is the + // app's core function, not "processing". It must NOT run in no-transcode + // (tag-only) mode: the resample stage was skipped there, so workingPath is + // still the user's ORIGINAL input file — these calls would destroy it (or + // fail outright on non-WAV input). Tag-only gets its ReplayGain figures + // from the read-only ebur128 measurement above and never alters audio. + if !n.noTranscode { + // LRA reduction toward the target range (self-gates if already under target). + const targetLRA = 7.0 + n.reduceLRA(workingPath, targetLRA, 4) + + // Normalize: measure the LRA-reduced signal and apply the linear offset to hit + // the loudness target. True peak is the conformance limiter's job below, so the + // full offset goes on here regardless of where it lands the peaks. + n.logFile.Write(fmt.Sprintf("Measuring loudness for %s", workingPath)) + lufs, _, lra, err := n.LUFS(workingPath) + if err != nil { + n.logFile.Write(fmt.Sprintf("error: %v", err)) + return false + } + n.logFile.Write(fmt.Sprintf("lra: %.1f", lra)) + lOffset := target - lufs + if err := n.Gain(workingPath, lOffset); err != nil { + n.appLog.Write(fmt.Sprintf("✗ Failed to apply gain: %s", filepath.Base(inputPath))) + n.logFile.Write(fmt.Sprintf("Gain application failed: %v", err)) + return false + } - // Normalize: measure the LRA-reduced signal and apply the linear offset to hit - // the loudness target. True peak is the conformance limiter's job below, so the - // full offset goes on here regardless of where it lands the peaks. - n.logFile.Write(fmt.Sprintf("Measuring loudness for %s", workingPath)) - lufs, _, lra, err := n.LUFS(workingPath) - if err != nil { - n.logFile.Write(fmt.Sprintf("error: %v", err)) - return false - } - n.logFile.Write(fmt.Sprintf("lra: %.1f", lra)) - lOffset := target - lufs - if err := n.Gain(workingPath, lOffset); err != nil { - n.appLog.Write(fmt.Sprintf("✗ Failed to apply gain: %s", filepath.Base(inputPath))) - n.logFile.Write(fmt.Sprintf("Gain application failed: %v", err)) - return false - } + // Encoder-specific character pre-limiter — lossy only. It pre-conditions the + // peaks with the headroom each codec/bitrate needs (softLimiterCalibration via + // strength) so the encoder's overshoot doesn't blow past the ceiling. PCM/FLAC + // have no encoder to overshoot, so they skip it. + if actualCodec != "PCM" && actualCodec != "flac" { + if err := n.CharacterLimit(workingPath, strength); err != nil { + n.logFile.Write(fmt.Sprintf("error in the encoder pre-limiter: %v", err)) + return false + } + } - // Encoder-specific character pre-limiter — lossy only. It pre-conditions the - // peaks with the headroom each codec/bitrate needs (softLimiterCalibration via - // strength) so the encoder's overshoot doesn't blow past the ceiling. PCM/FLAC - // have no encoder to overshoot, so they skip it. - if actualCodec != "PCM" && actualCodec != "flac" { - if err := n.CharacterLimit(workingPath, strength); err != nil { - n.logFile.Write(fmt.Sprintf("error in the encoder pre-limiter: %v", err)) + // Conformance limiter — true-peak, at the target TP ceiling. The + // final guarantee the output stays under the ceiling. + if err := n.conformLimit(workingPath, targetTp); err != nil { return false } } - // Conformance limiter — ALWAYS, true-peak, at the target TP ceiling. The - // final guarantee the output stays under the ceiling. - if err := n.conformLimit(workingPath, targetTp); err != nil { - return false - } - /* var loudnormFilterChain string if cfg.UseLoudnorm && measured != nil { @@ -1950,14 +1948,15 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool } */ - if _, isLossy := tiers[actualCodec]; isLossy { + // No filter stages in no-transcode mode: -af cannot combine with -c copy. + if _, isLossy := brightnessTiers[actualCodec]; isLossy && !n.noTranscode { if bf := n.buildBrightnessReduceFilterForLossy(strength); bf != "" { filterStages = append(filterStages, bf) } } // Add dithering for 16-bit PCM output. - if actualCodec == "PCM" && cfg.BitDepth == "16" { + if actualCodec == "PCM" && cfg.BitDepth == "16" && !n.noTranscode { filterStages = append(filterStages, "aresample=resampler=soxr:dither_method=high_shibata") } @@ -2006,14 +2005,18 @@ func (n *AudioNormalizer) processFile(inputPath string, cfg ProcessConfig) bool ) } - rho, err := n.PCMFileCoherence(workingPath) - if err != nil { - print(fmt.Errorf("there was an error in coherence check: %v", err)) - } - if rho > 0.4 { - print(fmt.Sprintf("Channel coherence is probably fine, value: %.1f", rho)) - } else { - print(fmt.Sprintf("Coherency might benefit from your attention: %.1f", rho)) + // Coherence check reads workingPath as WAV — meaningless (and failing, for + // non-WAV originals) in no-transcode mode where nothing was processed. + if !n.noTranscode { + rho, err := n.PCMFileCoherence(workingPath) + if err != nil { + print(fmt.Errorf("there was an error in coherence check: %v", err)) + } + if rho > 0.4 { + print(fmt.Sprintf("Channel coherence is probably fine, value: %.1f", rho)) + } else { + print(fmt.Sprintf("Coherency might benefit from your attention: %.1f", rho)) + } } n.logFile.Write("") diff --git a/go/notranscode_test.go b/go/notranscode_test.go new file mode 100644 index 0000000..3d03983 --- /dev/null +++ b/go/notranscode_test.go @@ -0,0 +1,105 @@ +package main + +import ( + "bytes" + "encoding/binary" + "math" + "os" + "path/filepath" + "testing" + + "github.com/fremen-fi/tnt/go/internal/ffmpeg" +) + +// writeToneWAV writes a 16-bit stereo PCM WAV with a ~-20 dBFS 440 Hz sine — +// the minimal valid input for the tag-only pipeline (ebur128 measurement plus +// a -c copy render). +func writeToneWAV(t *testing.T, path string, sampleRate int, seconds float64) { + t.Helper() + const channels = 2 + numSamples := int(float64(sampleRate) * seconds) + dataSize := numSamples * channels * 2 + + var buf bytes.Buffer + w := func(v any) { binary.Write(&buf, binary.LittleEndian, v) } + buf.WriteString("RIFF") + w(uint32(36 + dataSize)) + buf.WriteString("WAVE") + buf.WriteString("fmt ") + w(uint32(16)) + w(uint16(1)) // PCM + w(uint16(channels)) + w(uint32(sampleRate)) + w(uint32(sampleRate * channels * 2)) + w(uint16(channels * 2)) + w(uint16(16)) + buf.WriteString("data") + w(uint32(dataSize)) + amp := 0.1 * 32767.0 + for i := 0; i < numSamples; i++ { + s := int16(amp * math.Sin(2*math.Pi*440*float64(i)/float64(sampleRate))) + w(s) + w(s) + } + if err := os.WriteFile(path, buf.Bytes(), 0644); err != nil { + t.Fatal(err) + } +} + +// Regression test: in no-transcode (tag-only) mode processFile used to run +// the native loudness chain (reduceLRA / Gain / CharacterLimit / conformLimit) +// on workingPath, which — with the resample stage skipped — was still the +// user's ORIGINAL input file, silently rewriting it in place as a 32-bit +// float WAV (and failing outright on non-WAV input). Tag-only mode must never +// alter the input audio: the output is a -c copy of the source with +// ReplayGain tags from a read-only ebur128 measurement. +func TestProcessFileNoTranscodeLeavesOriginalUntouched(t *testing.T) { + if testing.Short() { + t.Skip("spawns ffmpeg") + } + ffmpegPath = ffmpeg.Path + + tmp := t.TempDir() + src := filepath.Join(tmp, "src.wav") + writeToneWAV(t, src, 44100, 1.0) + + orig, err := os.ReadFile(src) + if err != nil { + t.Fatal(err) + } + + outDir := filepath.Join(tmp, "out") + if err := os.MkdirAll(outDir, 0755); err != nil { + t.Fatal(err) + } + + n := &AudioNormalizer{ + logFile: &LogIntoFile{}, + appLog: &LogApp{}, + outputDir: outDir, + } + n.normalizationStandard = EBU + cfg := ProcessConfig{ + WriteTags: true, + NoTranscode: true, + Format: "MPEG-II L3", + } + n.applyConfig(cfg) + + if !n.processFile(src, cfg) { + t.Fatal("processFile returned false in tag-only (no-transcode) mode") + } + + after, err := os.ReadFile(src) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(orig, after) { + t.Fatalf("tag-only mode modified the original input file (size %d -> %d)", len(orig), len(after)) + } + + tagged := filepath.Join(outDir, "src.tagged.wav") + if _, err := os.Stat(tagged); err != nil { + t.Fatalf("expected tagged output at %s: %v", tagged, err) + } +} diff --git a/go/reduce_brightness.go b/go/reduce_brightness.go index 48b84e6..504bb97 100644 --- a/go/reduce_brightness.go +++ b/go/reduce_brightness.go @@ -6,6 +6,40 @@ import ( "github.com/fremen-fi/tnt/go/audio" ) +// brightnessTiers maps each lossy codec to the bitrate thresholds (kbps) at +// which brightness reduction and the encoder pre-limiter step up in strength. +// Codecs absent from the map (PCM, flac, native aac) never engage either stage. +type bitrateTiers struct { + mild, mid, hard int +} + +var brightnessTiers = map[string]bitrateTiers{ + "libmp3lame": {mild: 192, mid: 160, hard: 128}, + "aac_at": {mild: 128, mid: 96, hard: 64}, + "libfdk_aac": {mild: 128, mid: 96, hard: 64}, + "libopus": {mild: 96, mid: 64, hard: 48}, +} + +// brightnessStrength returns the brightness-reduction / encoder pre-limiter +// strength (0–3) for a codec at the given bitrate. The bitrate MUST be in +// kbps: processFile carries full bps for the needsFullNumber codecs, and +// comparing bps against these kbps tiers silently disabled the whole stage. +func brightnessStrength(codec string, bitrateKbps int) int { + t, ok := brightnessTiers[codec] + if !ok { + return 0 + } + switch { + case bitrateKbps <= t.hard: + return 3 + case bitrateKbps <= t.mid: + return 2 + case bitrateKbps <= t.mild: + return 1 + } + return 0 +} + func (n *AudioNormalizer) buildBrightnessReduceFilterForLossy(strength int) string { n.logFile.Write("Reducing brightness because of a lossy codec.") n.logFile.Write(fmt.Sprintf("Strength set to %d", strength)) diff --git a/go/reduce_brightness_test.go b/go/reduce_brightness_test.go new file mode 100644 index 0000000..638f1c7 --- /dev/null +++ b/go/reduce_brightness_test.go @@ -0,0 +1,58 @@ +package main + +import "testing" + +func TestBrightnessStrength(t *testing.T) { + cases := []struct { + codec string + bitrateKbps int + want int + }{ + {"libmp3lame", 320, 0}, + {"libmp3lame", 192, 1}, + {"libmp3lame", 160, 2}, + {"libmp3lame", 128, 3}, + {"libmp3lame", 96, 3}, + + {"libfdk_aac", 192, 0}, + {"libfdk_aac", 128, 1}, + {"libfdk_aac", 96, 2}, + {"libfdk_aac", 64, 3}, + + {"aac_at", 256, 0}, + {"aac_at", 128, 1}, + {"aac_at", 96, 2}, + {"aac_at", 64, 3}, + + {"libopus", 128, 0}, + {"libopus", 96, 1}, + {"libopus", 64, 2}, + {"libopus", 48, 3}, + + // Codecs without tiers never engage brightness reduction. + {"PCM", 128, 0}, + {"flac", 128, 0}, + {"aac", 64, 0}, + } + for _, c := range cases { + if got := brightnessStrength(c.codec, c.bitrateKbps); got != c.want { + t.Errorf("brightnessStrength(%q, %d) = %d, want %d", c.codec, c.bitrateKbps, got, c.want) + } + } +} + +// Regression: processFile holds the bitrate in full bps for the +// needsFullNumber codecs (e.g. "128k" → 128000). The tier comparison used to +// run on that bps value against kbps tiers, so strength was always 0 and the +// brightness reduction / pre-limiter calibration never engaged for those +// codecs. The fix converts back to kbps before selecting the tier. +func TestBrightnessStrengthUsesKbpsNotBps(t *testing.T) { + const bps = 128000 + if got := brightnessStrength("libmp3lame", bps/1000); got != 3 { + t.Errorf("libmp3lame at 128k: strength = %d, want 3", got) + } + // The old bug: passing bps yields 0 because 128000 > every kbps tier. + if got := brightnessStrength("libmp3lame", bps); got != 0 { + t.Errorf("sanity: bps input should miss all tiers, got %d", got) + } +} diff --git a/normalizer/README.md b/normalizer/README.md new file mode 100644 index 0000000..b9057c6 --- /dev/null +++ b/normalizer/README.md @@ -0,0 +1,87 @@ +# `normalizer` — streaming DSP module + +`github.com/fremen-fi/tnt/normalizer` — a self-contained, stdlib-only Go module +(own `go.mod`, no Wails, no GUI deps) holding the **streaming** version of TNT's +loudness + dynamics chain. Built for the **podcore** podcast transcoder, which +imports it (via a `replace` to this checkout; the TNT repo is public, so it can +also be imported by tag once pushed). + +Everything is `io.Reader → io.Writer`: the whole audio file is never resident. +The processors are causal one-pole followers (`compressor`, `upwardCompressor`) +plus a bounded-lookahead true-peak limiter (`conformanceLimiter`), so processing +memory is **O(lookahead)** — tens of KiB at any duration. ffmpeg does no loudness +or dynamics work; it only decodes the source to a raw f32 PCM file and encodes +the result. + +Public surface: `MeasureChain` (measure a PCM stream), `DeriveStage` (size one +adaptive stage from a measurement), `ApplyStage` (stream one stage), `Conform` +(gain-to-target + true-peak limit), plus `ChainMeasurement` / `StageParams`. + +podcore drives a **converging multi-pass** with these, over a working PCM file on +disk (so RAM stays low while arbitrarily large PCM lives on disk): + +``` +decode → file ; measure +repeat: DeriveStage(measurement) → ApplyStage(file→file') → measure(file') + until LRA ≤ target / diminishing returns / maxPasses +Conform(gain to -18 LUFS, true-peak limit to -5 dBTP) → file ; encode → AAC +``` + +Each pass re-reads the previous working file — the prior stages are already baked +in, so there is **no re-decode and nothing is AAC-encoded until the end**. On +Cloud Run the working file lives on a **gen2 ephemeral-disk volume** (not the +memory-backed filesystem), so the transcoder fits **1 vCPU / 1 GiB** at any +duration. + +## ⚠️ This duplicates the DSP math in `../go/audio` — by design, with a catch + +The per-sample math here is a **faithful port of the in-memory cores in +`go/audio`** (`processDynamics`, `UpwardCompressSamples`, `LookaheadLimitSamples`, +the BS.1770 meters, the `tpFIR` coefficients, the soft-knee gain curves). So the +same algorithms now live in **two places**: + +- `go/audio` — operates on resident `[]float64` buffers. The **desktop (Wails) + app** needs this: its editor/preview does random access over the whole file. +- `normalizer/` — operates on a stream. The **server** needs this: RAM stays + bounded (the multi-GiB working PCM lives on the ephemeral disk, streamed a + chunk at a time, never resident). + +This is *not* the old cross-repo duplication (podcore once vendored a byte-copy +of `go/audio` at `internal/audio` — that's deleted). It's two implementations of +the same math in one repo, for two genuinely different access patterns. + +**The catch — silent drift.** `normalizer/dynamics_test.go` verifies each +streaming processor against a *transcribed* whole-buffer reference (the `ref*` +functions in that test), **not** against `go/audio` directly. So if someone +changes a coefficient, threshold, or smoothing rule in `go/audio` and does **not** +mirror it here (and in the test's reference), the two engines diverge and nothing +fails. **When you touch the DSP math in `go/audio`, change it here too** (and the +test reference), or the desktop and server outputs will quietly differ. + +## Path to a single source of truth (deferred) + +The clean fix is to make `normalizer/` the *canonical* DSP and have `go/audio`'s +`*Samples` functions become thin wrappers that feed the whole buffer through the +streaming processors (a streaming processor run over an in-memory slice gives the +identical result — the streaming tests already prove that). Then there is one +implementation, and the in-memory API is just a convenience adapter. + +That's a desktop-app refactor (its DSP call sites assume in-place `[]float64` +mutation), deferred for now. Until it happens, treat the two as a mirrored pair. + +## Tuning note (podcore chain) + +`DeriveStage` sizes ONE stage from a measurement, mirroring `go/audio.reduceLRAPass`. +The downward Compressor and Upward Compressor **share** the LRA reduction (lower +peaks, lift quiet passages), so the range shrinks from both ends with far less +gain reduction than downward alone — downward-only buys ~1-2 LU and costs +dynamics; the pair does much more. + +This is the **real converging multi-pass**, the same feedback loop as the +in-memory `ReduceLRA`: each pass re-measures the working file and re-derives the +stage from the *current* (flattening) signal, so thresholds adapt and the loop +stops once it reaches the target or a pass buys < 0.3 LU. It lands precisely +(test: 15.6 LU → 6.8 LU in 3 passes, vs the earlier single-pass approximation +that stalled at 9.4). It costs extra disk passes, not re-decodes — the working +file is re-read, prior stages already baked in. Thresholds are a starting point; +dial them in against real spoken-word episodes. diff --git a/normalizer/dynamics.go b/normalizer/dynamics.go new file mode 100644 index 0000000..c6a2469 --- /dev/null +++ b/normalizer/dynamics.go @@ -0,0 +1,731 @@ +package normalizer + +// Streaming time-domain dynamics — the podcast loudness/dynamics chain that runs +// in the podcore transcoder without ever holding the whole file in memory. +// +// Every processor here is a faithful streaming port of the in-memory cores in +// the TNT desktop package (github.com/fremen-fi/tnt go/audio): same per-sample +// math, fed from an io.Reader instead of a resident []float64. The processors +// are all causal one-pole followers except the conformance limiter, which uses +// a bounded lookahead ring buffer — so memory is O(lookahead), a few tens of KiB +// regardless of file length. ffmpeg is only ever used by the caller to decode +// the source into this stream and to AAC-encode what comes out; it performs no +// loudness or dynamics work. +// +// The chain the podcore transcoder runs: +// +// measure → (downward) Compressor → Character Limiter → gain to target LUFS → +// Conformance Limiter (true-peak brick wall) +// +// The Compressor and Character Limiter share one feed-forward engine (the limiter +// is a compressor with an infinite ratio). The Conformance Limiter is a genuine +// true-peak limiter: its detection runs the BS.1770 4× polyphase FIR, so the +// reconstructed inter-sample peak is guaranteed under the ceiling — no oversampled +// ffmpeg alimiter needed. + +import ( + "io" + "math" +) + +// ── dB / smoothing helpers (ported verbatim from go/audio) ────────────────── + +// dbToLin converts decibels to linear amplitude. +func dbToLin(db float64) float64 { return math.Pow(10, db/20) } + +// linToDb converts linear amplitude to decibels; non-positive input floors at +// -100 dB (matching go/audio.LinearToDb — the dynamics math depends on this +// floor, NOT on -Inf, so silent samples produce a finite level). +func linToDb(lin float64) float64 { + if lin <= 0 { + return -100.0 + } + return 20 * math.Log10(lin) +} + +// onePoleCoef returns the one-pole smoothing coefficient for a timeMs time +// constant at sampleRate. A non-positive timeMs yields 0 (instantaneous move). +func onePoleCoef(timeMs, sampleRate float64) float64 { + if timeMs <= 0 { + return 0 + } + return math.Exp(-1.0 / (timeMs / 1000.0 * sampleRate)) +} + +// clampMag hard-limits x to the symmetric range [-m, m]. +func clampMag(x, m float64) float64 { + if x > m { + return m + } + if x < -m { + return -m + } + return x +} + +func clampf(x, lo, hi float64) float64 { return math.Max(lo, math.Min(hi, x)) } + +// staticGainDb returns the gain (dB, <= 0) the static soft-knee transfer curve +// applies to an input at levelDb, given threshold, ratio and a knee kneeDb wide +// centred on the threshold. ratio == +Inf makes it a hard limiter. Ported +// verbatim from go/audio.staticGainDb (Reiss & McPherson soft knee). +func staticGainDb(levelDb, thresholdDb, ratio, kneeDb float64) float64 { + slope := 1.0/ratio - 1.0 + over := levelDb - thresholdDb + switch { + case kneeDb > 0 && 2*over < -kneeDb: + return 0 + case kneeDb > 0 && 2*math.Abs(over) <= kneeDb: + x := over + kneeDb/2 + return slope * x * x / (2 * kneeDb) + default: + if over <= 0 { + return 0 + } + return slope * over + } +} + +// tpFIR is the order-48, 4-phase polyphase FIR from ITU-R BS.1770-5 Annex 2 +// (4× over-sampling for true-peak estimation). Ported verbatim from go/audio. +// Output = sum over k of tpFIR[phase][k]·x[n-k]; tap 0 is the most recent sample. +var tpFIR = [4][12]float64{ + { + 0.0017089843750, 0.0109863281250, -0.0196533203125, 0.0332031250000, + -0.0594482421875, 0.1373291015625, 0.9721679687500, -0.1022949218750, + 0.0476074218750, -0.0266113281250, 0.0148925781250, -0.0083007812500, + }, + { + -0.0291748046875, 0.0292968750000, -0.0517578125000, 0.0891113281250, + -0.1665039062500, 0.4650878906250, 0.7797851562500, -0.2003173828125, + 0.1015625000000, -0.0582275390625, 0.0330810546875, -0.0189208984375, + }, + { + -0.0189208984375, 0.0330810546875, -0.0582275390625, 0.1015625000000, + -0.2003173828125, 0.7797851562500, 0.4650878906250, -0.1665039062500, + 0.0891113281250, -0.0517578125000, 0.0292968750000, -0.0291748046875, + }, + { + -0.0083007812500, 0.0148925781250, -0.0266113281250, 0.0476074218750, + -0.1022949218750, 0.9721679687500, 0.1373291015625, -0.0594482421875, + 0.0332031250000, -0.0196533203125, 0.0109863281250, 0.0017089843750, + }, +} + +// tpHist is a 12-sample ring of one channel's recent input, the history the +// true-peak FIR reads back over. Samples older than the start of stream (or than +// 12 ago) read as absent. +type tpHist struct { + buf [12]float64 + count int +} + +func (h *tpHist) push(x float64) { + h.buf[h.count%12] = x + h.count++ +} + +// at returns the sample k positions back from the most recent (k==0 is current), +// and false if that sample predates the stream. +func (h *tpHist) at(k int) (float64, bool) { + idx := h.count - 1 - k + if idx < 0 { + return 0, false + } + return h.buf[idx%12], true +} + +// interpMax returns the maximum absolute inter-sample (true) peak at the current +// position across both channels — max over the 4 polyphase sub-samples of +// max(|FIR(L)|, |FIR(R)|). Stereo-linked, matching go/audio.requiredGainRange. +func interpMax(hL, hR *tpHist) float64 { + var mx float64 + for phase := 0; phase < 4; phase++ { + var accL, accR float64 + for k := 0; k < 12; k++ { + l, ok := hL.at(k) + if !ok { + break + } + r, _ := hR.at(k) + accL += tpFIR[phase][k] * l + accR += tpFIR[phase][k] * r + } + if a := math.Abs(accL); a > mx { + mx = a + } + if a := math.Abs(accR); a > mx { + mx = a + } + } + return mx +} + +// ── Compressor / Character Limiter (shared feed-forward engine) ───────────── + +// compressor is the streaming form of go/audio.processDynamics: a feed-forward +// compressor (finite ratio) or limiter (ratio == +Inf). Detection is peak +// (rmsMs <= 0) or one-pole-smoothed RMS (rmsMs > 0); the gain attacks when more +// reduction is wanted and releases when less is. Stereo-linked. O(1) state. +type compressor struct { + attCoef, relCoef float64 + rmsCoef float64 + useRMS bool + makeup float64 + thresholdDb float64 + ratio float64 + kneeDb float64 + + meanSq float64 + gain float64 +} + +func newCompressor(sampleRate, thresholdDb, ratio, kneeDb, attackMs, releaseMs, makeupDb, rmsMs float64) *compressor { + return &compressor{ + attCoef: onePoleCoef(attackMs, sampleRate), + relCoef: onePoleCoef(releaseMs, sampleRate), + rmsCoef: onePoleCoef(rmsMs, sampleRate), + useRMS: rmsMs > 0, + makeup: dbToLin(makeupDb), + thresholdDb: thresholdDb, + ratio: ratio, + kneeDb: kneeDb, + gain: 1.0, + } +} + +func (c *compressor) step(l, r float64) (float64, float64) { + var level float64 + if c.useRMS { + power := 0.5 * (l*l + r*r) + c.meanSq = c.rmsCoef*c.meanSq + (1-c.rmsCoef)*power + level = math.Sqrt(c.meanSq) + } else { + level = math.Max(math.Abs(l), math.Abs(r)) + } + target := dbToLin(staticGainDb(linToDb(level), c.thresholdDb, c.ratio, c.kneeDb)) + if target < c.gain { + c.gain = c.attCoef*c.gain + (1-c.attCoef)*target + } else { + c.gain = c.relCoef*c.gain + (1-c.relCoef)*target + } + return l * c.gain * c.makeup, r * c.gain * c.makeup +} + +// upwardGainDb returns the boost (>= 0 dB) an upward compressor applies to a +// passage at levelDb sitting BELOW thresholdDb: (threshold - level)·(1 - 1/ratio), +// soft-kneed and capped at maxBoostDb. Levels at or above the threshold get 0. +// The mirror of staticGainDb. Ported verbatim from go/audio.upwardGainDb. +func upwardGainDb(levelDb, thresholdDb, ratio, kneeDb, maxBoostDb float64) float64 { + slope := 1.0 - 1.0/ratio + under := thresholdDb - levelDb + var boost float64 + switch { + case kneeDb > 0 && 2*under < -kneeDb: + return 0 + case kneeDb > 0 && 2*math.Abs(under) <= kneeDb: + x := under + kneeDb/2 + boost = slope * x * x / (2 * kneeDb) + default: + if under <= 0 { + return 0 + } + boost = slope * under + } + if boost > maxBoostDb { + boost = maxBoostDb + } + return boost +} + +// upwardCompressor is the streaming form of go/audio.UpwardCompressSamples: it +// lifts passages below thresholdDb toward it, shrinking the loudness range from +// the bottom — the complement of the downward compressor, which lowers passages +// above threshold. Sharing the LRA reduction between the two needs far less gain +// reduction on the loud side, so fewer downward-compression artefacts and a much +// larger total LRA reduction than downward alone. Boost-only (gain >= 1), capped +// at maxBoostDb so quiet gaps / noise floor aren't dragged up without bound. +// O(1) state. +type upwardCompressor struct { + attCoef, relCoef float64 + rmsCoef float64 + useRMS bool + thresholdDb float64 + ratio float64 + kneeDb float64 + maxBoostDb float64 + + meanSq float64 + gain float64 +} + +func newUpwardCompressor(sampleRate, thresholdDb, ratio, kneeDb, attackMs, releaseMs, maxBoostDb, rmsMs float64) *upwardCompressor { + return &upwardCompressor{ + attCoef: onePoleCoef(attackMs, sampleRate), + relCoef: onePoleCoef(releaseMs, sampleRate), + rmsCoef: onePoleCoef(rmsMs, sampleRate), + useRMS: rmsMs > 0, + thresholdDb: thresholdDb, + ratio: ratio, + kneeDb: kneeDb, + maxBoostDb: maxBoostDb, + gain: 1.0, + } +} + +func (u *upwardCompressor) step(l, r float64) (float64, float64) { + var level float64 + if u.useRMS { + power := 0.5 * (l*l + r*r) + u.meanSq = u.rmsCoef*u.meanSq + (1-u.rmsCoef)*power + level = math.Sqrt(u.meanSq) + } else { + level = math.Max(math.Abs(l), math.Abs(r)) + } + target := dbToLin(upwardGainDb(linToDb(level), u.thresholdDb, u.ratio, u.kneeDb, u.maxBoostDb)) // >= 1 + if target > u.gain { + u.gain = u.attCoef*u.gain + (1-u.attCoef)*target // boost ramping in (attack) + } else { + u.gain = u.relCoef*u.gain + (1-u.relCoef)*target // backing off (release) + } + return l * u.gain, r * u.gain +} + +// ── Conformance limiter (streaming true-peak brick wall) ──────────────────── + +// conformanceLimiter is the streaming form of go/audio.LookaheadLimitSamples. +// For each sample it computes the gain needed to tame the true (inter-sample) +// peak to the ceiling, takes a sliding-window minimum over the lookahead window +// (so the gain is already down before a peak emerges — the conformance +// guarantee), smooths with instant attack / one-pole release, applies it and +// clamps at the ceiling as a brick-wall backstop. Output is delayed by `look` +// samples; Flush emits the tail. Memory is O(look). +type conformanceLimiter struct { + thresh float64 + look int + relCoef float64 + + histL, histR tpHist + + pendL, pendR []float64 // frames pushed but not yet emitted (FIFO) + dqIdx []int // monotonic deque: indices, values increasing from front + dqVal []float64 + + gain float64 + inIdx int + emitIdx int +} + +func newConformanceLimiter(sampleRate, thresholdDb, lookaheadMs, releaseMs float64) *conformanceLimiter { + look := int(math.Round(lookaheadMs / 1000.0 * sampleRate)) + if look < 1 { + look = 1 + } + return &conformanceLimiter{ + thresh: dbToLin(thresholdDb), + look: look, + relCoef: onePoleCoef(releaseMs, sampleRate), + gain: 1.0, + } +} + +// push feeds one input frame and returns one output frame once the lookahead +// window has filled (ok == false during the initial `look`-sample warm-up). +func (c *conformanceLimiter) push(l, r float64) (outL, outR float64, ok bool) { + i := c.inIdx + c.histL.push(l) + c.histR.push(r) + mx := interpMax(&c.histL, &c.histR) + req := 1.0 + if mx > c.thresh { + req = c.thresh / mx + } + + c.pendL = append(c.pendL, l) + c.pendR = append(c.pendR, r) + for len(c.dqVal) > 0 && c.dqVal[len(c.dqVal)-1] >= req { + c.dqVal = c.dqVal[:len(c.dqVal)-1] + c.dqIdx = c.dqIdx[:len(c.dqIdx)-1] + } + c.dqVal = append(c.dqVal, req) + c.dqIdx = append(c.dqIdx, i) + c.inIdx++ + + if i-c.emitIdx >= c.look { + outL, outR = c.emitFront() + return outL, outR, true + } + return 0, 0, false +} + +// emitFront emits the oldest buffered frame using the windowed-minimum gain. +func (c *conformanceLimiter) emitFront() (float64, float64) { + winMin := c.dqVal[0] + if winMin < c.gain { + c.gain = winMin // instant attack — the conformance guarantee + } else { + c.gain = c.relCoef*c.gain + (1-c.relCoef)*winMin + } + oL := clampMag(c.pendL[0]*c.gain, c.thresh) + oR := clampMag(c.pendR[0]*c.gain, c.thresh) + c.pendL = c.pendL[1:] + c.pendR = c.pendR[1:] + if c.dqIdx[0] == c.emitIdx { + c.dqVal = c.dqVal[1:] + c.dqIdx = c.dqIdx[1:] + } + c.emitIdx++ + return oL, oR +} + +// flush emits the remaining buffered frames (their lookahead windows truncate at +// end-of-stream), calling emit for each. +func (c *conformanceLimiter) flush(emit func(l, r float64)) { + for c.emitIdx < c.inIdx { + for len(c.dqIdx) > 0 && c.dqIdx[0] < c.emitIdx { + c.dqVal = c.dqVal[1:] + c.dqIdx = c.dqIdx[1:] + } + oL, oR := c.emitFront() + emit(oL, oR) + } +} + +// ── Measurement (single streaming pass) ───────────────────────────────────── + +// ChainMeasurement carries everything the chain needs to set its parameters: the +// EBU R128 figures plus the Dynamics-Score statistics (go/audio convention). +type ChainMeasurement struct { + IntegratedLUFS float64 // BS.1770 integrated loudness (channel-summed), -Inf for silence + LRA float64 // EBU R128 loudness range (LU) + TruePeakDB float64 // inter-sample true peak (dBTP), -Inf for silence + RMSLevelDB float64 // overall RMS level of channel 1 (dB) + RMSPeakDB float64 // windowed RMS-peak of channel 1 (dB) + CrestFactor float64 // peak / RMS of channel 1 + DynamicsScore float64 // sqrt(crest)·(RMSpeak−RMSlevel) +} + +const dsWindowSeconds = 0.05 // DS RMS-peak EMA time constant (astats default) + +type measurer struct { + filters [2]channelFilter + step int + stStep int + + stepSumSqL, stepSumSqR float64 + stepN int + stSumSqL, stSumSqR float64 + stN int + stepE, stStepE []float64 + + tpL, tpR tpHist + maxTP float64 + + dsAlpha float64 + dsSumSq, dsPeak, dsEMA, dsMaxE float64 + dsN int +} + +func newMeasurer(sampleRate int) *measurer { + fs := float64(sampleRate) + return &measurer{ + filters: newKweightFilters(fs), + step: sampleRate / 10, + stStep: sampleRate, + stepE: make([]float64, 0, 4096), + stStepE: make([]float64, 0, 4096), + dsAlpha: 1.0 - math.Exp(-1.0/(dsWindowSeconds*fs)), + } +} + +func (m *measurer) add(l, r float64) { + // K-weighted, channel-summed block energies for LUFS / LRA (BS.1770). + kl := m.filters[0].process(l) + kr := m.filters[1].process(r) + klSq, krSq := kl*kl, kr*kr + m.stepSumSqL += klSq + m.stepSumSqR += krSq + m.stepN++ + m.stSumSqL += klSq + m.stSumSqR += krSq + m.stN++ + if m.stepN >= m.step { + m.stepE = append(m.stepE, (m.stepSumSqL+m.stepSumSqR)/float64(m.stepN)) + m.stepSumSqL, m.stepSumSqR, m.stepN = 0, 0, 0 + } + if m.stN >= m.stStep { + m.stStepE = append(m.stStepE, (m.stSumSqL+m.stSumSqR)/float64(m.stN)) + m.stSumSqL, m.stSumSqR, m.stN = 0, 0, 0 + } + + // True peak across both channels. + m.tpL.push(l) + m.tpR.push(r) + if mx := interpMax(&m.tpL, &m.tpR); mx > m.maxTP { + m.maxTP = mx + } + + // Dynamics Score on channel 1 (left). + sq := l * l + m.dsSumSq += sq + if a := math.Abs(l); a > m.dsPeak { + m.dsPeak = a + } + m.dsEMA += m.dsAlpha * (sq - m.dsEMA) + if m.dsEMA > m.dsMaxE { + m.dsMaxE = m.dsEMA + } + m.dsN++ +} + +func (m *measurer) result() ChainMeasurement { + stepE := m.stepE + stStepE := m.stStepE + if m.stepN > 0 { + stepE = append(stepE, (m.stepSumSqL+m.stepSumSqR)/float64(m.stepN)) + } + if m.stN > 0 { + stStepE = append(stStepE, (m.stSumSqL+m.stSumSqR)/float64(m.stN)) + } + + res := ChainMeasurement{ + IntegratedLUFS: integratedLoudness(stepE, 4), // 400 ms block / 100 ms step + LRA: loudnessRange(stStepE, 3), // 3 s block / 1 s step + TruePeakDB: math.Inf(-1), + RMSLevelDB: -100, + RMSPeakDB: -100, + } + if m.maxTP > 0 { + res.TruePeakDB = 20 * math.Log10(m.maxTP) + } + if m.dsN > 0 { + rms := math.Sqrt(m.dsSumSq / float64(m.dsN)) + if rms > 0 { + res.RMSLevelDB = 20 * math.Log10(rms) + if m.dsMaxE > 0 { + res.RMSPeakDB = 10 * math.Log10(m.dsMaxE) + } + res.CrestFactor = m.dsPeak / rms + res.DynamicsScore = math.Sqrt(res.CrestFactor) * (res.RMSPeakDB - res.RMSLevelDB) + } + } + return res +} + +// ── Chain parameters & staged drivers ─────────────────────────────────────── + +// StageParams holds one gentle compressor→upward→character stage's settings +// (sample-rate-independent) — the streaming equivalent of one go/audio +// reduceLRAPass. The converging multi-pass re-derives these each pass from the +// freshly re-measured signal, so the thresholds track the flattening signal and +// the stage auto-tapers toward the target. +type StageParams struct { + CompThresholdDb float64 + CompRatio float64 + CompKneeDb float64 + CompAttackMs float64 + CompReleaseMs float64 + CompRMSWindowMs float64 + + UpThresholdDb float64 + UpRatio float64 + UpKneeDb float64 + UpAttackMs float64 + UpReleaseMs float64 + UpMaxBoostDb float64 + UpRMSWindowMs float64 + + CharThresholdDb float64 + CharKneeDb float64 + CharAttackMs float64 + CharReleaseMs float64 +} + +// compModifiers reproduces go/audio.GetCompressionModifiers: DS-driven attack / +// release / ratio multipliers. +func compModifiers(ds float64) (attack, release, ratioMul float64) { + switch { + case ds < 9.0: + return 4.0, 4.0, 0.15 + case ds < 15.0: + return 2.0, 2.0, 2.1 + case ds <= 21.0: + return 1.0, 1.0, 1.0 + default: + excess := math.Min(ds-21.0, 79.0) + scale := 2.0 + (2.0 * excess / 79.0) + return 1.0 / scale, 1.0 / scale, 4.0 + (4.0 * excess / 79.0) + } +} + +// DeriveStage sizes ONE compressor→upward→character stage from a measurement, +// mirroring go/audio.reduceLRAPass. The downward compressor lowers loud passages +// (threshold between the RMS level and peak, ratio sized to the overshoot); the +// upward compressor lifts quiet ones toward the RMS level, so the pair shrinks +// LRA from both ends with far less gain reduction than downward alone; the slow +// character limiter catches the residual. The converging loop calls this fresh +// on each re-measured pass, so it adapts and tapers toward the target. +func DeriveStage(m ChainMeasurement, targetLRA float64) StageParams { + const wideKnee = 12.0 + const rmsWindowMs = 300.0 + const charKnee = 6.0 + + over := m.LRA - targetLRA + if over < 0 { + over = 0 + } + att, rel, _ := compModifiers(m.DynamicsScore) + span := m.RMSPeakDB - m.RMSLevelDB + dyn := clampf((m.DynamicsScore-9)/12, 0, 1) + interp := clampf(0.6-over*0.05, 0.1, 0.6) + + return StageParams{ + CompThresholdDb: m.RMSLevelDB + interp*span, + CompRatio: clampf(1+over*0.15, 1, 6), + CompKneeDb: wideKnee, + CompAttackMs: 100 * att, + CompReleaseMs: 1500 * rel, + CompRMSWindowMs: rmsWindowMs, + + UpThresholdDb: m.RMSLevelDB, + UpRatio: clampf(1+over*0.1, 1, 4), + UpKneeDb: wideKnee, + UpAttackMs: 300 * att, + UpReleaseMs: 600 * rel, + UpMaxBoostDb: 5.0, + UpRMSWindowMs: rmsWindowMs, + + CharThresholdDb: m.RMSPeakDB - (0.2+0.3*dyn)*span, + CharKneeDb: charKnee, + CharAttackMs: 150 * att, + CharReleaseMs: 1200 * rel, + } +} + +func (sp StageParams) build(sampleRate int) (comp *compressor, up *upwardCompressor, char *compressor) { + fs := float64(sampleRate) + comp = newCompressor(fs, sp.CompThresholdDb, sp.CompRatio, sp.CompKneeDb, + sp.CompAttackMs, sp.CompReleaseMs, 0, sp.CompRMSWindowMs) + up = newUpwardCompressor(fs, sp.UpThresholdDb, sp.UpRatio, sp.UpKneeDb, + sp.UpAttackMs, sp.UpReleaseMs, sp.UpMaxBoostDb, sp.UpRMSWindowMs) + // Character limiter: infinite-ratio soft-knee, peak detector. + char = newCompressor(fs, sp.CharThresholdDb, math.Inf(1), sp.CharKneeDb, + sp.CharAttackMs, sp.CharReleaseMs, 0, 0) + return comp, up, char +} + +// streamFrames reads interleaved f32 stereo from src, applies fn to each frame, +// and writes the result to dst. For zero-latency per-frame processors. +func streamFrames(src io.Reader, dst io.Writer, fn func(l, r float64) (float64, float64)) error { + buf := make([]float32, frameChunk*2) + for { + n, rerr := readF32LE(src, buf) + for i := 0; i+1 < n; i += 2 { + l, r := fn(float64(buf[i]), float64(buf[i+1])) + buf[i], buf[i+1] = float32(l), float32(r) + } + if n > 0 { + if werr := writeF32LE(dst, buf[:n]); werr != nil { + return werr + } + } + if rerr == io.EOF || rerr == io.ErrUnexpectedEOF { + return nil + } + if rerr != nil { + return rerr + } + } +} + +// ApplyStage streams src→dst applying one compressor→upward→character stage. +// The converging multi-pass calls this once per pass, re-deriving sp from the +// freshly measured intermediate each time. +func ApplyStage(src io.Reader, dst io.Writer, sampleRate int, sp StageParams) error { + comp, up, char := sp.build(sampleRate) + return streamFrames(src, dst, func(l, r float64) (float64, float64) { + l, r = comp.step(l, r) + l, r = up.step(l, r) + l, r = char.step(l, r) + return l, r + }) +} + +// Conform streams src→dst applying a constant gainDB then the true-peak +// conformance limiter at ceilingDB — the final pass after loudness/LRA +// processing. lookaheadMs / releaseMs shape the limiter; memory is O(lookahead). +func Conform(src io.Reader, dst io.Writer, sampleRate int, gainDB, ceilingDB, lookaheadMs, releaseMs float64) error { + gain := dbToLin(gainDB) + lim := newConformanceLimiter(float64(sampleRate), ceilingDB, lookaheadMs, releaseMs) + + in := make([]float32, frameChunk*2) + out := make([]float32, 0, frameChunk*2) + flushOut := func() error { + if len(out) == 0 { + return nil + } + if err := writeF32LE(dst, out); err != nil { + return err + } + out = out[:0] + return nil + } + for { + n, rerr := readF32LE(src, in) + for i := 0; i+1 < n; i += 2 { + l := float64(in[i]) * gain + r := float64(in[i+1]) * gain + if ol, or, ok := lim.push(l, r); ok { + out = append(out, float32(ol), float32(or)) + if len(out) >= cap(out) { + if err := flushOut(); err != nil { + return err + } + } + } + } + if rerr == io.EOF || rerr == io.ErrUnexpectedEOF { + break + } + if rerr != nil { + return rerr + } + } + var flushErr error + lim.flush(func(l, r float64) { + out = append(out, float32(l), float32(r)) + if len(out) >= cap(out) { + if err := flushOut(); err != nil && flushErr == nil { + flushErr = err + } + } + }) + if flushErr != nil { + return flushErr + } + return flushOut() +} + +// ── Streaming drivers ─────────────────────────────────────────────────────── + +const frameChunk = 8192 // stereo frames processed per read + +// MeasureChain measures src (interleaved f32 stereo PCM) in a single streaming +// pass. Read to EOF; src need not be seekable. +func MeasureChain(src io.Reader, sampleRate int) (ChainMeasurement, error) { + m := newMeasurer(sampleRate) + buf := make([]float32, frameChunk*2) + for { + n, err := readF32LE(src, buf) + for i := 0; i+1 < n; i += 2 { + m.add(float64(buf[i]), float64(buf[i+1])) + } + if err == io.EOF || err == io.ErrUnexpectedEOF { + break + } + if err != nil { + return ChainMeasurement{}, err + } + } + return m.result(), nil +} diff --git a/normalizer/dynamics_test.go b/normalizer/dynamics_test.go new file mode 100644 index 0000000..069f79b --- /dev/null +++ b/normalizer/dynamics_test.go @@ -0,0 +1,426 @@ +package normalizer + +import ( + "bytes" + "math" + "testing" +) + +// ── reference implementations (straightforward whole-buffer ports of the +// go/audio cores) the streaming processors must match sample-for-sample ─────── + +func refProcessDynamics(left, right []float64, sr, thresholdDb, ratio, kneeDb, attackMs, releaseMs, makeupDb, rmsMs float64) { + n := min(len(left), len(right)) + attCoef := onePoleCoef(attackMs, sr) + relCoef := onePoleCoef(releaseMs, sr) + rmsCoef := onePoleCoef(rmsMs, sr) + makeup := dbToLin(makeupDb) + useRMS := rmsMs > 0 + var meanSq, gain float64 + gain = 1.0 + for i := 0; i < n; i++ { + var level float64 + if useRMS { + power := 0.5 * (left[i]*left[i] + right[i]*right[i]) + meanSq = rmsCoef*meanSq + (1-rmsCoef)*power + level = math.Sqrt(meanSq) + } else { + level = math.Max(math.Abs(left[i]), math.Abs(right[i])) + } + target := dbToLin(staticGainDb(linToDb(level), thresholdDb, ratio, kneeDb)) + if target < gain { + gain = attCoef*gain + (1-attCoef)*target + } else { + gain = relCoef*gain + (1-relCoef)*target + } + left[i] *= gain * makeup + right[i] *= gain * makeup + } +} + +func refLookaheadLimit(left, right []float64, sr, thresholdDb, lookaheadMs, releaseMs float64) { + n := min(len(left), len(right)) + if n == 0 { + return + } + thresh := dbToLin(thresholdDb) + look := int(math.Round(lookaheadMs / 1000.0 * sr)) + if look < 1 { + look = 1 + } + if look > n-1 { + look = n - 1 + } + required := make([]float64, n) + for i := 0; i < n; i++ { + var mx float64 + for phase := 0; phase < 4; phase++ { + var accL, accR float64 + for k := 0; k < 12; k++ { + idx := i - k + if idx < 0 { + break + } + accL += tpFIR[phase][k] * left[idx] + accR += tpFIR[phase][k] * right[idx] + } + if a := math.Abs(accL); a > mx { + mx = a + } + if a := math.Abs(accR); a > mx { + mx = a + } + } + if mx > thresh { + required[i] = thresh / mx + } else { + required[i] = 1 + } + } + // forward sliding-window minimum, window [i, i+look] + winMin := make([]float64, n) + for i := 0; i < n; i++ { + m := required[i] + end := min(i+look, n-1) + for j := i + 1; j <= end; j++ { + if required[j] < m { + m = required[j] + } + } + winMin[i] = m + } + relCoef := onePoleCoef(releaseMs, sr) + gain := 1.0 + for i := 0; i < n; i++ { + t := winMin[i] + if t < gain { + gain = t + } else { + gain = relCoef*gain + (1-relCoef)*t + } + left[i] = clampMag(left[i]*gain, thresh) + right[i] = clampMag(right[i]*gain, thresh) + } +} + +// testSignal builds a deterministic stereo signal: mixed sines plus periodic +// transient spikes (so the limiter has real true-peak work to do). +func testSignal(n int) (l, r []float64) { + l = make([]float64, n) + r = make([]float64, n) + for i := 0; i < n; i++ { + t := float64(i) + v := 0.4*math.Sin(2*math.Pi*440*t/44100) + 0.3*math.Sin(2*math.Pi*1300*t/44100) + if i%2000 == 0 { + v += 0.9 // transient + } + l[i] = v + r[i] = 0.95 * v + } + return l, r +} + +func maxAbsDiff(a, b []float64) float64 { + var d float64 + for i := range a { + if x := math.Abs(a[i] - b[i]); x > d { + d = x + } + } + return d +} + +func TestCompressorMatchesReference(t *testing.T) { + const sr = 44100 + l, r := testSignal(50000) + refL, refR := append([]float64(nil), l...), append([]float64(nil), r...) + refProcessDynamics(refL, refR, sr, -24, 3, 6, 100, 1500, 0, 300) + + c := newCompressor(sr, -24, 3, 6, 100, 1500, 0, 300) + outL := make([]float64, len(l)) + outR := make([]float64, len(r)) + for i := range l { + outL[i], outR[i] = c.step(l[i], r[i]) + } + if d := maxAbsDiff(outL, refL); d > 1e-12 { + t.Errorf("compressor L diverges from reference by %g", d) + } + if d := maxAbsDiff(outR, refR); d > 1e-12 { + t.Errorf("compressor R diverges from reference by %g", d) + } +} + +func TestCharacterLimiterMatchesReference(t *testing.T) { + const sr = 44100 + l, r := testSignal(40000) + refL, refR := append([]float64(nil), l...), append([]float64(nil), r...) + refProcessDynamics(refL, refR, sr, -10, math.Inf(1), 6, 150, 1200, 0, 0) + + c := newCompressor(sr, -10, math.Inf(1), 6, 150, 1200, 0, 0) + outL := make([]float64, len(l)) + outR := make([]float64, len(r)) + for i := range l { + outL[i], outR[i] = c.step(l[i], r[i]) + } + if d := maxAbsDiff(outL, refL); d > 1e-12 { + t.Errorf("character limiter L diverges by %g", d) + } + if d := maxAbsDiff(outR, refR); d > 1e-12 { + t.Errorf("character limiter R diverges by %g", d) + } +} + +func refUpwardCompress(left, right []float64, sr, thresholdDb, ratio, kneeDb, attackMs, releaseMs, maxBoostDb, rmsMs float64) { + n := min(len(left), len(right)) + attCoef := onePoleCoef(attackMs, sr) + relCoef := onePoleCoef(releaseMs, sr) + rmsCoef := onePoleCoef(rmsMs, sr) + useRMS := rmsMs > 0 + var meanSq float64 + gain := 1.0 + for i := 0; i < n; i++ { + var level float64 + if useRMS { + power := 0.5 * (left[i]*left[i] + right[i]*right[i]) + meanSq = rmsCoef*meanSq + (1-rmsCoef)*power + level = math.Sqrt(meanSq) + } else { + level = math.Max(math.Abs(left[i]), math.Abs(right[i])) + } + target := dbToLin(upwardGainDb(linToDb(level), thresholdDb, ratio, kneeDb, maxBoostDb)) + if target > gain { + gain = attCoef*gain + (1-attCoef)*target + } else { + gain = relCoef*gain + (1-relCoef)*target + } + left[i] *= gain + right[i] *= gain + } +} + +func TestUpwardCompressorMatchesReference(t *testing.T) { + const sr = 44100 + l, r := testSignal(50000) + for i := range l { // scale so material sits both sides of the threshold + l[i] *= 0.3 + r[i] *= 0.3 + } + refL, refR := append([]float64(nil), l...), append([]float64(nil), r...) + refUpwardCompress(refL, refR, sr, -24, 3, 12, 300, 600, 5, 300) + + u := newUpwardCompressor(sr, -24, 3, 12, 300, 600, 5, 300) + outL := make([]float64, len(l)) + outR := make([]float64, len(r)) + for i := range l { + outL[i], outR[i] = u.step(l[i], r[i]) + } + if d := maxAbsDiff(outL, refL); d > 1e-12 { + t.Errorf("upward compressor L diverges by %g", d) + } + if d := maxAbsDiff(outR, refR); d > 1e-12 { + t.Errorf("upward compressor R diverges by %g", d) + } +} + +// convergeChain runs the converging multi-pass over in-memory PCM, mirroring the +// podcore orchestration: measure → derive one stage → apply → re-measure → stop +// at target / diminishing returns / maxPasses, then gain-to-target + conformance +// limit. Returns the final output's measurement, the output PCM, and pass count. +func convergeChain(t *testing.T, sr int, pcm []byte, targetLUFS, targetTP, targetLRA float64) (ChainMeasurement, []byte, int) { + t.Helper() + const maxPasses = 6 + m, err := MeasureChain(bytes.NewReader(pcm), sr) + if err != nil { + t.Fatal(err) + } + cur := pcm + prevLRA := m.LRA + passes := 0 + for passes < maxPasses && !math.IsInf(m.IntegratedLUFS, -1) && m.LRA > targetLRA { + sp := DeriveStage(m, targetLRA) + var next bytes.Buffer + if err := ApplyStage(bytes.NewReader(cur), &next, sr, sp); err != nil { + t.Fatal(err) + } + cur = next.Bytes() + if m, err = MeasureChain(bytes.NewReader(cur), sr); err != nil { + t.Fatal(err) + } + passes++ + if prevLRA-m.LRA < 0.3 { // diminishing returns + break + } + prevLRA = m.LRA + } + gainDB := 0.0 + if !math.IsInf(m.IntegratedLUFS, -1) { + gainDB = targetLUFS - m.IntegratedLUFS + } + var out bytes.Buffer + if err := Conform(bytes.NewReader(cur), &out, sr, gainDB, targetTP, 20, 80); err != nil { + t.Fatal(err) + } + fm, err := MeasureChain(bytes.NewReader(out.Bytes()), sr) + if err != nil { + t.Fatal(err) + } + return fm, out.Bytes(), passes +} + +// TestChainReducesLRA confirms the converging multi-pass meaningfully shrinks the +// loudness range — downward compression alone buys only ~0.5-1 LU. +func TestChainReducesLRA(t *testing.T) { + const sr = 44100 + const seg = 3 * sr // 3 s segments + n := seg * 6 // alternating loud / quiet, 18 s + buf := make([]float32, n*2) + for i := 0; i < n; i++ { + amp := 0.3 // loud + if (i/seg)%2 == 1 { + amp = 0.05 // quiet + } + v := float32(amp * math.Sin(2*math.Pi*300*float64(i)/sr)) + buf[i*2] = v + buf[i*2+1] = v + } + var src bytes.Buffer + writeF32LE(&src, buf) + + m1, err := MeasureChain(bytes.NewReader(src.Bytes()), sr) + if err != nil { + t.Fatal(err) + } + final, _, passes := convergeChain(t, sr, src.Bytes(), -18, -5, 7) + t.Logf("LRA: input=%.1f LU → %.1f LU in %d passes (Δ %.1f)", m1.LRA, final.LRA, passes, m1.LRA-final.LRA) + if final.LRA >= m1.LRA { + t.Errorf("chain did not reduce LRA: %.1f → %.1f", m1.LRA, final.LRA) + } + if m1.LRA-final.LRA < 3 { + t.Errorf("LRA reduction too small: %.1f → %.1f (Δ %.1f LU)", m1.LRA, final.LRA, m1.LRA-final.LRA) + } +} + +func TestConformanceLimiterMatchesReference(t *testing.T) { + const sr = 44100 + l, r := testSignal(60000) + refL, refR := append([]float64(nil), l...), append([]float64(nil), r...) + refLookaheadLimit(refL, refR, sr, -5, 20, 80) + + lim := newConformanceLimiter(sr, -5, 20, 80) + outL := make([]float64, 0, len(l)) + outR := make([]float64, 0, len(r)) + for i := range l { + if ol, or, ok := lim.push(l[i], r[i]); ok { + outL = append(outL, ol) + outR = append(outR, or) + } + } + lim.flush(func(a, b float64) { + outL = append(outL, a) + outR = append(outR, b) + }) + + if len(outL) != len(l) { + t.Fatalf("streaming limiter emitted %d frames, want %d", len(outL), len(l)) + } + if d := maxAbsDiff(outL, refL); d > 1e-12 { + t.Errorf("conformance limiter L diverges by %g", d) + } + if d := maxAbsDiff(outR, refR); d > 1e-12 { + t.Errorf("conformance limiter R diverges by %g", d) + } +} + +// TestConformanceLimiterGuarantee: the streaming limiter's output true peak must +// not exceed the ceiling (the whole point of a conformance limiter). +func TestConformanceLimiterGuarantee(t *testing.T) { + const sr = 44100 + const ceilingDB = -5.0 + l, r := testSignal(60000) + lim := newConformanceLimiter(sr, ceilingDB, 20, 80) + var outL, outR []float64 + for i := range l { + if ol, or, ok := lim.push(l[i], r[i]); ok { + outL = append(outL, ol) + outR = append(outR, or) + } + } + lim.flush(func(a, b float64) { outL = append(outL, a); outR = append(outR, b) }) + + // Measure true peak of the output the same way the meter does. + var hL, hR tpHist + var maxTP float64 + for i := range outL { + hL.push(outL[i]) + hR.push(outR[i]) + if mx := interpMax(&hL, &hR); mx > maxTP { + maxTP = mx + } + } + ceiling := dbToLin(ceilingDB) + if maxTP > ceiling*(1+1e-6) { + t.Errorf("output true peak %.4f dBTP exceeds ceiling %.1f dBTP", 20*math.Log10(maxTP), ceilingDB) + } +} + +func TestMeasureChainSilence(t *testing.T) { + var buf bytes.Buffer + writeF32LE(&buf, make([]float32, 44100*2*2)) // 2 s stereo silence + m, err := MeasureChain(&buf, 44100) + if err != nil { + t.Fatal(err) + } + if !math.IsInf(m.IntegratedLUFS, -1) { + t.Errorf("silence integrated LUFS = %v, want -Inf", m.IntegratedLUFS) + } +} + +func TestMeasureChainSineLoudness(t *testing.T) { + // −23 dBFS 997 Hz sine, channel-summed BS.1770: a stereo (dual-mono) tone + // reads ~3 LU above the single-channel figure, and K-weighting near 1 kHz is + // close to flat, so expect roughly −20 LUFS. Generous bounds. + const sr = 44100 + amp := math.Pow(10, -23.0/20) + n := sr * 8 + buf := make([]float32, n*2) + for i := 0; i < n; i++ { + v := float32(amp * math.Sin(2*math.Pi*997*float64(i)/sr)) + buf[i*2] = v + buf[i*2+1] = v + } + var b bytes.Buffer + writeF32LE(&b, buf) + m, err := MeasureChain(&b, sr) + if err != nil { + t.Fatal(err) + } + if m.IntegratedLUFS < -24 || m.IntegratedLUFS > -16 { + t.Errorf("sine integrated LUFS = %.2f, want roughly [-24,-16]", m.IntegratedLUFS) + } +} + +// TestConvergeEndToEnd: a quiet, dynamic signal through the converging chain +// lands near the target LUFS and under the true-peak ceiling. +func TestConvergeEndToEnd(t *testing.T) { + const sr = 44100 + l, r := testSignal(sr * 6) + for i := range l { // scale down so it needs positive gain + l[i] *= 0.2 + r[i] *= 0.2 + } + inter := make([]float32, len(l)*2) + for i := range l { + inter[i*2] = float32(l[i]) + inter[i*2+1] = float32(r[i]) + } + var srcBuf bytes.Buffer + writeF32LE(&srcBuf, inter) + + final, _, _ := convergeChain(t, sr, srcBuf.Bytes(), -18, -5, 7) + if math.Abs(final.IntegratedLUFS-(-18.0)) > 1.5 { + t.Errorf("final LUFS = %.2f, want within 1.5 of -18", final.IntegratedLUFS) + } + if final.TruePeakDB > -5.0+0.1 { + t.Errorf("final true peak = %.2f dBTP, exceeds -5 ceiling", final.TruePeakDB) + } +} diff --git a/normalizer/gain.go b/normalizer/gain.go new file mode 100644 index 0000000..f946a05 --- /dev/null +++ b/normalizer/gain.go @@ -0,0 +1,59 @@ +package normalizer + +// Streaming gain application and sample-peak limiter. +// +// The limiter here is a sample-peak (not true-peak) limiter. For the podcore +// transcode pipeline the encode pass applies FFmpeg loudnorm+alimiter which +// enforces the TP ceiling; this stage only prevents gross over-shoots before +// encoding and adds negligible CPU cost compared to the FIR-based TP approach. + +import ( + "io" + "math" +) + +// ApplyGainAndLimit reads interleaved f32 stereo samples from src, multiplies +// each sample by gainLinear, then hard-clips at |peakCeiling| (linear, e.g. +// 0.5623 for −5 dBFS), and writes the result to dst. +// +// Returns the number of sample frames processed and the highest |sample| seen +// BEFORE the limiter (to let the caller detect how much headroom was consumed). +func ApplyGainAndLimit(src io.Reader, dst io.Writer, gainLinear, peakCeiling float64) (frames int64, preGainPeakDB float64, err error) { + buf := make([]float32, 8192) // 4096 stereo frames per chunk + var maxAbs float64 + + for { + n, rerr := readF32LE(src, buf) + if n > 0 { + for i := range n { + s := float64(buf[i]) * gainLinear + if a := math.Abs(s); a > maxAbs { + maxAbs = a + } + if s > peakCeiling { + s = peakCeiling + } else if s < -peakCeiling { + s = -peakCeiling + } + buf[i] = float32(s) + } + if werr := writeF32LE(dst, buf[:n]); werr != nil { + return frames, 0, werr + } + frames += int64(n / 2) // stereo: 2 samples per frame + } + if rerr == io.EOF || rerr == io.ErrUnexpectedEOF { + break + } + if rerr != nil { + return frames, 0, rerr + } + } + + if maxAbs > 0 { + preGainPeakDB = 20 * math.Log10(maxAbs) + } else { + preGainPeakDB = -math.MaxFloat64 + } + return frames, preGainPeakDB, nil +} diff --git a/normalizer/go.mod b/normalizer/go.mod new file mode 100644 index 0000000..63974be --- /dev/null +++ b/normalizer/go.mod @@ -0,0 +1,3 @@ +module github.com/fremen-fi/tnt/normalizer + +go 1.22 diff --git a/normalizer/lufs.go b/normalizer/lufs.go new file mode 100644 index 0000000..56a4699 --- /dev/null +++ b/normalizer/lufs.go @@ -0,0 +1,344 @@ +package normalizer + +// Streaming EBU R128 loudness measurement (ITU-R BS.1770-4). +// +// Memory model: block energies are stored as one float64 per 100 ms step. +// A 3-hour stereo file at 44100 Hz produces at most 108 000 blocks ≈ 864 KiB. +// No sample arrays are allocated beyond the processing buffer passed in. + +import ( + "fmt" + "io" + "math" + "sort" +) + +// LoudnessResult carries the measurements returned by MeasureLUFS. +type LoudnessResult struct { + IntegratedLUFS float64 // EBU R128 integrated program loudness (LUFS) + LRA float64 // Loudness range (LU) + SamplePeakDB float64 // Highest sample peak across all channels (dBFS) +} + +// MeasureLUFS performs a single streaming pass over interleaved f32 stereo PCM +// read from r, computing EBU R128 integrated loudness, LRA, and sample peak. +// +// sampleRate must be the PCM stream's sample rate (e.g. 44100). +// r must be positioned at the first PCM sample byte (after the WAV data header). +// dataBytes is the number of PCM bytes to read (-1 = read until EOF). +func MeasureLUFS(r io.Reader, sampleRate int, dataBytes int64) (LoudnessResult, error) { + if sampleRate <= 0 { + return LoudnessResult{}, fmt.Errorf("invalid sample rate %d", sampleRate) + } + + filters := newKweightFilters(float64(sampleRate)) + + // Block sizes in samples (per channel). + step := sampleRate / 10 // 100 ms step + blockLen := 4 * sampleRate / 10 // 400 ms block = 4 steps + + // Short-term block sizes for LRA. + stStep := sampleRate // 1 s step + stBlockLen := 3 * sampleRate // 3 s block + + // We accumulate a ring buffer of the last blockLen per-step mean-squares + // so we can form the 400 ms sliding window without re-reading samples. + stepEnergies := make([]float64, 0, 4096) // mean-square per 100 ms step + stStepEnergies := make([]float64, 0, 4096) + + // Per-step accumulators (reset every `step` samples). + var stepSumSq float64 + var stepSamples int + + // Per-short-term-step (1 s) accumulators. + var stSumSq float64 + var stSamples int + + var samplePeak float64 // max |sample| (linear) + + // Read buffer: process in chunks of `step` interleaved samples (2 channels). + chunkSize := step * 2 // interleaved L+R + buf := make([]float32, chunkSize) + + var totalRead int64 + + for { + if dataBytes >= 0 && totalRead >= dataBytes { + break + } + toRead := chunkSize + if dataBytes >= 0 { + bytesLeft := dataBytes - totalRead + samplesLeft := int(bytesLeft / 4) + if samplesLeft < toRead { + toRead = samplesLeft + } + } + if toRead == 0 { + break + } + + n, err := readF32LE(r, buf[:toRead]) + totalRead += int64(n * 4) + + for i := 0; i < n; i += 2 { + l := float64(buf[i]) + r := float64(buf[i+1]) + + // Track sample peak before K-weighting. + if a := math.Abs(l); a > samplePeak { + samplePeak = a + } + if a := math.Abs(r); a > samplePeak { + samplePeak = a + } + + // Apply K-weighting per channel. + kl := filters[0].process(l) + kr := filters[1].process(r) + + // Mean square contribution (average over channels). + ms := (kl*kl + kr*kr) * 0.5 + stepSumSq += ms + stepSamples++ + + stSumSq += ms + stSamples++ + } + + // Emit a 100 ms block when full. + if stepSamples >= step { + stepEnergies = append(stepEnergies, stepSumSq/float64(stepSamples)) + stepSumSq = 0 + stepSamples = 0 + } + + // Emit a 1 s short-term block. + if stSamples >= stStep { + stStepEnergies = append(stStepEnergies, stSumSq/float64(stSamples)) + stSumSq = 0 + stSamples = 0 + } + + if err == io.EOF || err == io.ErrUnexpectedEOF { + break + } + if err != nil { + return LoudnessResult{}, err + } + } + + // Flush partial step. + if stepSamples > 0 { + stepEnergies = append(stepEnergies, stepSumSq/float64(stepSamples)) + } + if stSamples > 0 { + stStepEnergies = append(stStepEnergies, stSumSq/float64(stSamples)) + } + + intLUFS := integratedLoudness(stepEnergies, blockLen/step) + lra := loudnessRange(stStepEnergies, stBlockLen/stStep) + + var peakDB float64 + if samplePeak > 0 { + peakDB = 20 * math.Log10(samplePeak) + } else { + peakDB = -math.MaxFloat64 + } + + return LoudnessResult{ + IntegratedLUFS: intLUFS, + LRA: lra, + SamplePeakDB: peakDB, + }, nil +} + +// integratedLoudness computes EBU R128 integrated loudness from per-step mean +// squares. stepsPerBlock is the number of steps that make up one gating block +// (4 for 100 ms steps / 400 ms blocks). +func integratedLoudness(stepMSq []float64, stepsPerBlock int) float64 { + if len(stepMSq) < stepsPerBlock { + return math.Inf(-1) + } + + // Build 400 ms sliding-window block powers. + blocks := make([]float64, 0, len(stepMSq)) + for i := 0; i+stepsPerBlock <= len(stepMSq); i++ { + var sum float64 + for j := range stepsPerBlock { + sum += stepMSq[i+j] + } + blocks = append(blocks, sum/float64(stepsPerBlock)) + } + + // Absolute gate: -70 LUFS → linear threshold. + absThresh := math.Pow(10, (-70.0+0.691)/10) // 0.691 is the offset for mean-square + + // First pass: mean of blocks above absolute gate. + var sumAbove float64 + var countAbove int + for _, b := range blocks { + if b >= absThresh { + sumAbove += b + countAbove++ + } + } + if countAbove == 0 { + return math.Inf(-1) + } + ungatedMean := sumAbove / float64(countAbove) + + // Relative gate: ungated mean − 10 LU. + relThresh := ungatedMean * math.Pow(10, -10.0/10) + + // Second pass: mean of blocks above relative gate. + var sumRel float64 + var countRel int + for _, b := range blocks { + if b >= relThresh { + sumRel += b + countRel++ + } + } + if countRel == 0 { + return math.Inf(-1) + } + + integratedMean := sumRel / float64(countRel) + return -0.691 + 10*math.Log10(integratedMean) +} + +// loudnessRange computes LRA from per-1s-step mean squares. +// stepsPerSTBlock = 3 (3 s blocks with 1 s steps). +func loudnessRange(stStepMSq []float64, stepsPerSTBlock int) float64 { + if len(stStepMSq) < stepsPerSTBlock { + return 0 + } + + stBlocks := make([]float64, 0, len(stStepMSq)) + for i := 0; i+stepsPerSTBlock <= len(stStepMSq); i++ { + var sum float64 + for j := range stepsPerSTBlock { + sum += stStepMSq[i+j] + } + stBlocks = append(stBlocks, sum/float64(stepsPerSTBlock)) + } + + // Absolute gate: -70 LUFS. + absThresh := math.Pow(10, (-70.0+0.691)/10) + var above []float64 + for _, b := range stBlocks { + if b >= absThresh { + above = append(above, b) + } + } + if len(above) == 0 { + return 0 + } + + // Relative gate: mean of above-absolute minus 20 LU. + var sumA float64 + for _, b := range above { + sumA += b + } + relThresh := (sumA / float64(len(above))) * math.Pow(10, -20.0/10) + + var gated []float64 + for _, b := range above { + if b >= relThresh { + gated = append(gated, b) + } + } + if len(gated) < 2 { + return 0 + } + + // Sort and take 10th – 95th percentile range. + sort.Float64s(gated) + lo := gated[int(math.Round(float64(len(gated)-1)*0.10))] + hi := gated[int(math.Round(float64(len(gated)-1)*0.95))] + + var luLo, luHi float64 + if lo > 0 { + luLo = -0.691 + 10*math.Log10(lo) + } + if hi > 0 { + luHi = -0.691 + 10*math.Log10(hi) + } + lra := luHi - luLo + if lra < 0 { + return 0 + } + return lra +} + +// kweightFilter is a single biquad IIR section (Direct Form I). +type kweightFilter struct { + b0, b1, b2 float64 + a1, a2 float64 + x1, x2 float64 // input delay + y1, y2 float64 // output delay +} + +func (f *kweightFilter) process(x float64) float64 { + y := f.b0*x + f.b1*f.x1 + f.b2*f.x2 - f.a1*f.y1 - f.a2*f.y2 + f.x2, f.x1 = f.x1, x + f.y2, f.y1 = f.y1, y + return y +} + +// channelFilter chains stage1 (high-shelf pre-filter) and stage2 (high-pass). +type channelFilter struct { + s1, s2 kweightFilter +} + +func (cf *channelFilter) process(x float64) float64 { + return cf.s2.process(cf.s1.process(x)) +} + +// newKweightFilters returns K-weighting biquad pairs for stereo (L, R). +// Coefficients are derived analytically from ITU-R BS.1770-4 / libebur128. +func newKweightFilters(fs float64) [2]channelFilter { + var cf [2]channelFilter + for i := range cf { + cf[i].s1 = kweightStage1(fs) + cf[i].s2 = kweightStage2(fs) + } + return cf +} + +func kweightStage1(fs float64) kweightFilter { + // High-shelf pre-filter: f0=1681.97 Hz, G=+4 dB, Q=0.707 + f0 := 1681.974450955533 + G := 3.999843853973347 + Q := 0.7071752369554196 + + K := math.Tan(math.Pi * f0 / fs) + Vh := math.Pow(10.0, G/20.0) + Vb := math.Pow(Vh, 0.4845) + + a0 := 1.0 + K/Q + K*K + return kweightFilter{ + b0: (Vh + Vb*K/Q + K*K) / a0, + b1: 2.0 * (K*K - Vh) / a0, + b2: (Vh - Vb*K/Q + K*K) / a0, + a1: 2.0 * (K*K - 1.0) / a0, + a2: (1.0 - K/Q + K*K) / a0, + } +} + +func kweightStage2(fs float64) kweightFilter { + // High-pass (RLB) filter: f0=38.14 Hz, Q=0.500 + f0 := 38.13547087692325 + Q := 0.5003270373238773 + + K := math.Tan(math.Pi * f0 / fs) + a0 := 1.0 + K/Q + K*K + return kweightFilter{ + b0: 1.0 / a0, + b1: -2.0 / a0, + b2: 1.0 / a0, + a1: 2.0 * (K*K - 1.0) / a0, + a2: (1.0 - K/Q + K*K) / a0, + } +} diff --git a/normalizer/normalizer.go b/normalizer/normalizer.go new file mode 100644 index 0000000..b76b76a --- /dev/null +++ b/normalizer/normalizer.go @@ -0,0 +1,133 @@ +// Package normalizer implements streaming EBU R128 loudness normalization for +// PCM WAV audio. All processing uses io.Reader / io.Writer so the entire audio +// file is never held in memory; memory use is O(duration/100ms) for loudness +// gate block energies, typically < 1 MiB even for a 3-hour file. +// +// Intended for use with podcore's Cloud Run Job transcoder. The caller decodes +// audio to an f32le stereo WAV with FFmpeg, then calls Normalize which: +// +// 1. Reads the WAV header and validates the RIFF data-chunk size against the +// actual file size (guards against ffmpeg's silent 4 GiB RIFF wrap). +// 2. Performs a single streaming EBU R128 measurement pass to obtain +// integrated LUFS, LRA, and sample peak. +// 3. Computes the linear gain required to reach TargetLUFS. +// 4. Performs a second streaming pass applying that gain and hard-clipping at +// SamplePeakCeiling (not true-peak; the encode pass handles TP via loudnorm +// linear mode + alimiter). +// +// The output is a new f32le stereo WAV that the caller passes to the encode +// stage. +package normalizer + +import ( + "fmt" + "io" + "math" + "os" +) + +// Config holds normalization parameters. +type Config struct { + // TargetLUFS is the desired integrated loudness, e.g. -18.0. + TargetLUFS float64 + + // SamplePeakCeiling is the maximum allowed |sample| value (linear). + // Use math.Pow(10, dBTP/20) — e.g. 0.5623 for -5 dBFS. + SamplePeakCeiling float64 + + // MaxLRA, when > 0, clamps the measured LRA fed to the caller. When the + // input was pre-processed by a dynamics stage the caller should set this to + // the LRA target (e.g. 11) so that loudnorm linear mode is always selected. + // Pass 0 to return the raw measured value. + MaxLRA float64 +} + +// Result carries the measurements and actions taken by Normalize. +type Result struct { + MeasuredLUFS float64 + MeasuredLRA float64 // raw measured (before MaxLRA clamping) + ClampedLRA float64 // value fed back to encode stage (clamped if MaxLRA > 0) + MeasuredPeakDB float64 // sample peak before gain, dBFS + AppliedGainDB float64 + Clipped bool // any sample was at SamplePeakCeiling after gain +} + +// Normalize reads the WAV at srcPath, normalizes it to cfg.TargetLUFS and +// writes the result to dstPath (overwritten or created). Both paths should +// reference ephemeral disk locations to avoid RAM pressure. +// +// The function is intentionally two-pass (measure then apply) so it can run on +// arbitrarily long files with bounded memory. An ephemeral disk write for the +// output is the only I/O amplification. +func Normalize(srcPath, dstPath string, cfg Config) (Result, error) { + // ── Pass 1: measure ────────────────────────────────────────────────────── + src, err := os.Open(srcPath) + if err != nil { + return Result{}, fmt.Errorf("open src: %w", err) + } + defer src.Close() + + info, err := ReadWAVInfo(src) + if err != nil { + return Result{}, fmt.Errorf("read WAV header: %w", err) + } + if info.Channels != 2 { + return Result{}, fmt.Errorf("expected stereo (2-ch) WAV, got %d channels", info.Channels) + } + + meas, err := MeasureLUFS(src, info.SampleRate, info.DataSize) + if err != nil { + return Result{}, fmt.Errorf("measure LUFS: %w", err) + } + + gainDB := cfg.TargetLUFS - meas.IntegratedLUFS + gainLinear := math.Pow(10, gainDB/20) + + clampedLRA := meas.LRA + if cfg.MaxLRA > 0 && clampedLRA > cfg.MaxLRA { + clampedLRA = cfg.MaxLRA + } + + // ── Pass 2: apply gain + sample-peak limit ─────────────────────────────── + // Seek src back to start of PCM data. + if _, err := src.Seek(0, io.SeekStart); err != nil { + return Result{}, err + } + if _, err := ReadWAVInfo(src); err != nil { + return Result{}, fmt.Errorf("re-read header for pass 2: %w", err) + } + + dst, err := os.Create(dstPath) + if err != nil { + return Result{}, fmt.Errorf("create dst: %w", err) + } + defer dst.Close() + + if err := WriteWAVHeader(dst, info.SampleRate, info.Channels, info.DataSize); err != nil { + return Result{}, fmt.Errorf("write WAV header: %w", err) + } + + _, postGainPeakDB, err := ApplyGainAndLimit(src, dst, gainLinear, cfg.SamplePeakCeiling) + if err != nil { + return Result{}, fmt.Errorf("apply gain: %w", err) + } + + // postGainPeakDB is peak AFTER gain, before hard-clip. + clipped := postGainPeakDB > linearToDBFS(cfg.SamplePeakCeiling) + + return Result{ + MeasuredLUFS: meas.IntegratedLUFS, + MeasuredLRA: meas.LRA, + ClampedLRA: clampedLRA, + MeasuredPeakDB: meas.SamplePeakDB, + AppliedGainDB: gainDB, + Clipped: clipped, + }, nil +} + +func linearToDBFS(linear float64) float64 { + if linear <= 0 { + return math.Inf(-1) + } + return 20 * math.Log10(linear) +} diff --git a/normalizer/normalizer_test.go b/normalizer/normalizer_test.go new file mode 100644 index 0000000..b8740c6 --- /dev/null +++ b/normalizer/normalizer_test.go @@ -0,0 +1,159 @@ +package normalizer + +import ( + "bytes" + "encoding/binary" + "math" + "os" + "testing" +) + +// sineTone generates interleaved f32 stereo sine samples at 440 Hz. +func sineTone(sampleRate, durationSec int, amplitudeLinear float64) []float32 { + n := sampleRate * durationSec * 2 // stereo + buf := make([]float32, n) + for i := range sampleRate * durationSec { + v := float32(amplitudeLinear * math.Sin(2*math.Pi*440*float64(i)/float64(sampleRate))) + buf[i*2] = v + buf[i*2+1] = v + } + return buf +} + +func makeWAV(t *testing.T, sampleRate int, samples []float32) string { + t.Helper() + f, err := os.CreateTemp(t.TempDir(), "test-*.wav") + if err != nil { + t.Fatal(err) + } + dataBytes := int64(len(samples) * 4) + if err := WriteWAVHeader(f, sampleRate, 2, dataBytes); err != nil { + t.Fatal(err) + } + if err := writeF32LE(f, samples); err != nil { + t.Fatal(err) + } + f.Close() + return f.Name() +} + +func TestMeasureLUFS_SilenceIsNegInf(t *testing.T) { + samples := make([]float32, 44100*5*2) // 5 s silence + path := makeWAV(t, 44100, samples) + f, _ := os.Open(path) + defer f.Close() + info, _ := ReadWAVInfo(f) + res, err := MeasureLUFS(f, info.SampleRate, info.DataSize) + if err != nil { + t.Fatal(err) + } + if !math.IsInf(res.IntegratedLUFS, -1) { + t.Errorf("silence: expected -Inf LUFS, got %v", res.IntegratedLUFS) + } +} + +func TestMeasureLUFS_SineReasonable(t *testing.T) { + // 10 s of 440 Hz sine at -23 dBFS: LUFS should be around -23 ± 2. + amp := math.Pow(10, -23.0/20) + samples := sineTone(44100, 10, amp) + path := makeWAV(t, 44100, samples) + f, _ := os.Open(path) + defer f.Close() + info, _ := ReadWAVInfo(f) + res, err := MeasureLUFS(f, info.SampleRate, info.DataSize) + if err != nil { + t.Fatal(err) + } + // Sine has a known relationship between amplitude and LUFS after K-weighting. + // Tolerance ±3 LU covers the K-weighting attenuation at 440 Hz. + if res.IntegratedLUFS < -30 || res.IntegratedLUFS > -18 { + t.Errorf("LUFS %v out of expected range [-30, -18]", res.IntegratedLUFS) + } +} + +func TestNormalize_GainApplied(t *testing.T) { + amp := math.Pow(10, -30.0/20) // quiet signal at -30 dBFS amplitude + samples := sineTone(44100, 10, amp) + src := makeWAV(t, 44100, samples) + dst := src + ".out.wav" + t.Cleanup(func() { os.Remove(dst) }) + + cfg := Config{ + TargetLUFS: -18.0, + SamplePeakCeiling: math.Pow(10, -5.0/20), + MaxLRA: 11, + } + res, err := Normalize(src, dst, cfg) + if err != nil { + t.Fatal(err) + } + if res.AppliedGainDB <= 0 { + t.Errorf("expected positive gain for quiet signal, got %v dB", res.AppliedGainDB) + } + + // The output WAV should be readable. + outF, err := os.Open(dst) + if err != nil { + t.Fatal("output file not created:", err) + } + defer outF.Close() + outInfo, err := ReadWAVInfo(outF) + if err != nil { + t.Fatal("output WAV header invalid:", err) + } + if outInfo.SampleRate != 44100 || outInfo.Channels != 2 { + t.Errorf("unexpected output format: %+v", outInfo) + } +} + +func TestReadWAVInfo_TruncatedRejected(t *testing.T) { + // Build a WAV with a data chunk size at the RIFF limit. + var buf bytes.Buffer + buf.WriteString("RIFF") + binary.Write(&buf, binary.LittleEndian, uint32(0xFFFFFFFF)) + buf.WriteString("WAVE") + buf.WriteString("fmt ") + binary.Write(&buf, binary.LittleEndian, uint32(18)) + binary.Write(&buf, binary.LittleEndian, uint16(3)) // IEEE float + binary.Write(&buf, binary.LittleEndian, uint16(2)) // channels + binary.Write(&buf, binary.LittleEndian, uint32(44100)) + binary.Write(&buf, binary.LittleEndian, uint32(44100*8)) + binary.Write(&buf, binary.LittleEndian, uint16(8)) + binary.Write(&buf, binary.LittleEndian, uint16(32)) + binary.Write(&buf, binary.LittleEndian, uint16(0)) + buf.WriteString("data") + binary.Write(&buf, binary.LittleEndian, uint32(0xFFFFFFFF)) // clamped + + r := bytes.NewReader(buf.Bytes()) + _, err := ReadWAVInfo(r) + if err == nil { + t.Error("expected error for RIFF-limit data chunk, got nil") + } +} + +func TestApplyGainAndLimit_Clips(t *testing.T) { + // Signal at 0 dBFS; gain of +6 dB should exceed ceiling of -5 dBFS. + samples := []float32{0.9, 0.9, -0.9, -0.9} // 2 stereo frames + var src bytes.Buffer + writeF32LE(&src, samples) + + ceiling := math.Pow(10, -5.0/20) // 0.5623 + gain := math.Pow(10, 6.0/20) // 2.0× + var dst bytes.Buffer + _, peakDB, err := ApplyGainAndLimit(&src, &dst, gain, ceiling) + if err != nil { + t.Fatal(err) + } + // Post-gain peak should be > ceiling (signal clipped). + if peakDB <= 20*math.Log10(ceiling) { + t.Errorf("expected peak above ceiling, got %v dB", peakDB) + } + // Output samples should be clamped. + out := make([]float32, 4) + binary.Read(bytes.NewReader(dst.Bytes()), binary.LittleEndian, &out) + for _, s := range out { + if math.Abs(float64(s)) > ceiling+1e-6 { + t.Errorf("output sample %v exceeds ceiling %v", s, ceiling) + } + } +} diff --git a/normalizer/pcm.go b/normalizer/pcm.go new file mode 100644 index 0000000..605226a --- /dev/null +++ b/normalizer/pcm.go @@ -0,0 +1,196 @@ +package normalizer + +import ( + "encoding/binary" + "fmt" + "io" + "math" + "os" +) + +// riffSafeLimit is 4 GiB minus a small margin. A WAV dataSize at or above +// this value means ffmpeg clamped the RIFF uint32 — the file is truncated. +const riffSafeLimit = int64(0xFFFF0000) + +// WAVInfo describes the PCM stream found in a WAV file. +type WAVInfo struct { + SampleRate int + Channels int + Format uint16 // 1 = PCM int, 3 = IEEE float + BitDepth int + DataSize int64 // bytes of PCM data +} + +// ReadWAVInfo reads the RIFF/WAVE header from r (which must support Seek). +// On success r is positioned at the first PCM sample byte. +// Returns an error if the data-chunk size is at the RIFF 4 GiB ceiling +// (indicating ffmpeg truncated a file that exceeded the limit). +func ReadWAVInfo(r io.ReadSeeker) (WAVInfo, error) { + var hdr [4]byte + if _, err := io.ReadFull(r, hdr[:]); err != nil { + return WAVInfo{}, fmt.Errorf("read RIFF tag: %w", err) + } + if string(hdr[:]) != "RIFF" { + return WAVInfo{}, fmt.Errorf("not a WAV file (expected RIFF, got %q)", string(hdr[:])) + } + var riffSize uint32 + if err := binary.Read(r, binary.LittleEndian, &riffSize); err != nil { + return WAVInfo{}, err + } + if _, err := io.ReadFull(r, hdr[:]); err != nil { + return WAVInfo{}, fmt.Errorf("read WAVE tag: %w", err) + } + if string(hdr[:]) != "WAVE" { + return WAVInfo{}, fmt.Errorf("not a WAV file (expected WAVE, got %q)", string(hdr[:])) + } + + var info WAVInfo + for { + var chunkID [4]byte + if _, err := io.ReadFull(r, chunkID[:]); err != nil { + return WAVInfo{}, fmt.Errorf("read chunk id: %w", err) + } + var chunkSize uint32 + if err := binary.Read(r, binary.LittleEndian, &chunkSize); err != nil { + return WAVInfo{}, err + } + switch string(chunkID[:]) { + case "fmt ": + if err := parseFmtChunk(r, chunkSize, &info); err != nil { + return WAVInfo{}, err + } + case "data": + if int64(chunkSize) >= riffSafeLimit { + return WAVInfo{}, fmt.Errorf( + "WAV data chunk size 0x%X is at RIFF 4 GiB ceiling — file is truncated", + chunkSize, + ) + } + // Cross-check against actual file size when we have an *os.File. + if f, ok := r.(*os.File); ok { + stat, err := f.Stat() + if err == nil { + pos, _ := f.Seek(0, io.SeekCurrent) + remaining := stat.Size() - pos + if int64(chunkSize) > remaining+4 { + return WAVInfo{}, fmt.Errorf( + "WAV dataSize %d > actual remaining bytes %d — file truncated", + chunkSize, remaining, + ) + } + } + } + info.DataSize = int64(chunkSize) + return info, nil + default: + // Skip unknown or non-essential chunks (LIST, ID3, …). + if _, err := io.CopyN(io.Discard, r, int64(chunkSize)); err != nil { + return WAVInfo{}, err + } + if chunkSize%2 != 0 { + r.Seek(1, io.SeekCurrent) //nolint:errcheck + } + } + } +} + +func parseFmtChunk(r io.Reader, size uint32, info *WAVInfo) error { + var audioFmt uint16 + if err := binary.Read(r, binary.LittleEndian, &audioFmt); err != nil { + return err + } + var channels uint16 + if err := binary.Read(r, binary.LittleEndian, &channels); err != nil { + return err + } + var sampleRate uint32 + if err := binary.Read(r, binary.LittleEndian, &sampleRate); err != nil { + return err + } + io.CopyN(io.Discard, r, 4) // byte rate (ignored) + io.CopyN(io.Discard, r, 2) // block align (ignored) + var bitDepth uint16 + if err := binary.Read(r, binary.LittleEndian, &bitDepth); err != nil { + return err + } + // Skip any extra fmt bytes (extensible format, etc.) + read := uint32(16) + if size > read { + io.CopyN(io.Discard, r, int64(size-read)) //nolint:errcheck + } + info.Format = audioFmt + info.Channels = int(channels) + info.SampleRate = int(sampleRate) + info.BitDepth = int(bitDepth) + return nil +} + +// WriteWAVHeader writes a minimal RIFF/WAVE/fmt/data header for f32le stereo. +// dataSize is the number of PCM bytes that will follow; pass 0 to write a +// placeholder (the caller is responsible for seeking back and patching it). +func WriteWAVHeader(w io.Writer, sampleRate, channels int, dataSize int64) error { + blockAlign := channels * 4 // f32 = 4 bytes per sample per channel + byteRate := sampleRate * blockAlign + + riffSize := uint32(36 + dataSize) // 36 = fmt chunk (24) + data header (8) + + writes := []any{ + []byte("RIFF"), + riffSize, + []byte("WAVE"), + []byte("fmt "), + uint32(18), // fmt chunk size (PCM would be 16; f32 needs extra 2 bytes) + uint16(3), // IEEE float + uint16(channels), + uint32(sampleRate), + uint32(byteRate), + uint16(blockAlign), + uint16(32), // bit depth + uint16(0), // cbSize extra bytes (0 for non-extensible) + []byte("data"), + uint32(dataSize), + } + for _, v := range writes { + if err := binary.Write(w, binary.LittleEndian, v); err != nil { + return err + } + } + return nil +} + +// SamplesFromBytes returns the number of stereo frame pairs a WAVInfo's DataSize +// represents (DataSize / (channels * bytesPerSample)). +func (w WAVInfo) SampleFrames() int64 { + bytesPerSample := int64(w.BitDepth / 8) + if bytesPerSample == 0 { + return 0 + } + return w.DataSize / (int64(w.Channels) * bytesPerSample) +} + +// readF32LE reads up to len(buf) interleaved float32 samples from r. +func readF32LE(r io.Reader, buf []float32) (int, error) { + b := make([]byte, len(buf)*4) + n, err := io.ReadFull(r, b) + if err == io.ErrUnexpectedEOF { + // Partial read at end of stream — convert what we got. + n = n &^ 3 // round down to float32 boundary + err = io.EOF + } + samples := n / 4 + for i := range samples { + bits := binary.LittleEndian.Uint32(b[i*4:]) + buf[i] = math.Float32frombits(bits) + } + return samples, err +} + +// writeF32LE writes interleaved float32 samples to w. +func writeF32LE(w io.Writer, buf []float32) error { + b := make([]byte, len(buf)*4) + for i, s := range buf { + binary.LittleEndian.PutUint32(b[i*4:], math.Float32bits(s)) + } + _, err := w.Write(b) + return err +}