From 5ec82eb1fd41fb42144ec6091929b1db9060a9de Mon Sep 17 00:00:00 2001 From: Brandon McGinty Date: Thu, 20 Aug 2026 14:27:20 -0400 Subject: [PATCH] fix recorded audio truncation and rate mismatch Accumulate incoming frames per source and consume fixed-size chunks. Each speaker's frames were queued whole and one frame was taken per tick, so a frame that did not match the recorder's frame size was truncated or padded and the recording drifted out of time with the audio. Frames of any size now append to a per-source buffer that is drained in exact chunks. Tell the recorder whether a frame is stereo instead of guessing from its length. Mono microphone frames were being interpreted as stereo, which halved their duration and produced static in the output. Let the recorder worker close ffmpeg's stdin. Stop closed it from the caller while the worker was still writing, turning a normal stop into a broken pipe and losing the tail of the recording. Reserve the output file exclusively and hand ffmpeg the descriptor. The path was generated, then reopened by name, so another process could take the name in between. The file is opened once with O_EXCL and passed to ffmpeg as an inherited descriptor. Co-Authored-By: Claude Opus 5 --- gumble/gumbleopenal/stream.go | 8 +-- recording/recorder.go | 114 ++++++++++++++++++++++++---------- recording/recorder_test.go | 85 +++++++++++++++++++++++-- 3 files changed, 165 insertions(+), 42 deletions(-) diff --git a/gumble/gumbleopenal/stream.go b/gumble/gumbleopenal/stream.go index a3f16b4..dd502ec 100644 --- a/gumble/gumbleopenal/stream.go +++ b/gumble/gumbleopenal/stream.go @@ -26,7 +26,7 @@ type FilePlayer interface { } type Recorder interface { - RecordAudioFrame(source uint32, samples []int16) + RecordAudioFrame(source uint32, samples []int16, stereo bool) } const recorderOutgoingSource uint32 = ^uint32(0) @@ -351,7 +351,7 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) { } } if recorder != nil && recordPtr > 0 { - recorder.RecordAudioFrame(e.User.Session, recordBuffer[:recordPtr]) + recorder.RecordAudioFrame(e.User.Session, recordBuffer[:recordPtr], true) } reclaim() @@ -500,13 +500,13 @@ func (s *Stream) sourceRoutine(inputDevice *string) { // Send stereo buffer when file is playing outgoing <- gumble.AudioBuffer(outputBuffer) if recorder := s.getRecorder(); recorder != nil { - recorder.RecordAudioFrame(recorderOutgoingSource, outputBuffer) + recorder.RecordAudioFrame(recorderOutgoingSource, outputBuffer, true) } } else if hasMicInput { // Send mic when no file is playing outgoing <- gumble.AudioBuffer(int16Buffer) if recorder := s.getRecorder(); recorder != nil { - recorder.RecordAudioFrame(recorderOutgoingSource, int16Buffer) + recorder.RecordAudioFrame(recorderOutgoingSource, int16Buffer, false) } } } diff --git a/recording/recorder.go b/recording/recorder.go index 5346c01..82d896d 100644 --- a/recording/recorder.go +++ b/recording/recorder.go @@ -57,14 +57,30 @@ func New(directory string, format string, now time.Time, frameSize int, interval if err := os.MkdirAll(directory, 0755); err != nil { return nil, err } - path := UniquePath(directory, now, format) - args := ffmpegArgs(format, path) - cmd := exec.Command("ffmpeg", args...) - stdin, err := cmd.StdinPipe() + output, path, err := reserveOutput(directory, now, format) if err != nil { return nil, err } + args := ffmpegArgs(format) + cmd := exec.Command("ffmpeg", args...) + // Pass the reserved file descriptor directly to ffmpeg. The file is never + // reopened by pathname, preventing replacement between reservation and use. + cmd.ExtraFiles = []*os.File{output} + stdin, err := cmd.StdinPipe() + if err != nil { + _ = output.Close() + _ = os.Remove(path) + return nil, err + } if err := cmd.Start(); err != nil { + _ = output.Close() + _ = os.Remove(path) + return nil, err + } + if err := output.Close(); err != nil { + _ = cmd.Process.Kill() + _ = cmd.Wait() + _ = os.Remove(path) return nil, err } recorder := &Recorder{ @@ -91,6 +107,33 @@ func NormalizeFormat(format string) string { return format } +func reserveOutput(directory string, now time.Time, format string) (*os.File, string, error) { + for { + path := UniquePath(directory, now, format) + file, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600) + if os.IsExist(err) { + continue + } + if err != nil { + return nil, "", err + } + return file, path, nil + } +} + +// reservePath remains available for callers that only need to reserve a name. +func reservePath(directory string, now time.Time, format string) (string, error) { + file, path, err := reserveOutput(directory, now, format) + if err != nil { + return "", err + } + if err := file.Close(); err != nil { + _ = os.Remove(path) + return "", err + } + return path, nil +} + func UniquePath(directory string, now time.Time, format string) string { base := fmt.Sprintf("barnard-recording-%s", now.Format("20060102-150405")) path := filepath.Join(directory, base+"."+format) @@ -109,14 +152,14 @@ func (r *Recorder) Path() string { return r.path } -func (r *Recorder) RecordAudioFrame(source uint32, samples []int16) { +func (r *Recorder) RecordAudioFrame(source uint32, samples []int16, stereo bool) { if r == nil || len(samples) == 0 { return } if len(r.input) >= cap(r.input) { return } - frame := NormalizeStereoFrame(samples, r.frameSize) + frame := NormalizeStereoFrame(samples, stereo) select { case r.input <- sourceFrame{source: source, samples: frame}: default: @@ -128,10 +171,9 @@ func (r *Recorder) Stop() error { return nil } r.once.Do(func() { + // run owns stdin and closes it only after it has stopped writing. + // Closing it here races writePCM and turns a normal stop into EPIPE. close(r.stop) - if r.stdin != nil { - r.stdin.Close() - } }) select { case <-r.done: @@ -150,31 +192,35 @@ func (r *Recorder) run() { defer close(r.done) ticker := time.NewTicker(r.interval) defer ticker.Stop() - queues := make(map[uint32][][]int16) - frame := make([]int16, r.frameSize*gumble.AudioChannels) + // Per-source accumulated stereo samples. Incoming frames of any size are + // appended and then consumed in frameSize*AudioChannels chunks each tick. + queues := make(map[uint32][]int16) + chunkSize := r.frameSize * gumble.AudioChannels + chunk := make([]int16, chunkSize) for { select { case <-r.stop: r.closeEncoder() return case item := <-r.input: - queues[item.source] = append(queues[item.source], item.samples) + queues[item.source] = append(queues[item.source], item.samples...) case <-ticker.C: - clear(frame) - for source, queue := range queues { - if len(queue) == 0 { + clear(chunk) + for source, buffer := range queues { + if len(buffer) == 0 { delete(queues, source) continue } - mix(frame, queue[0]) - queue = queue[1:] - if len(queue) == 0 { + // Mix one chunk worth of samples from this source. + if len(buffer) <= chunkSize { + mix(chunk, buffer) delete(queues, source) } else { - queues[source] = queue + mix(chunk, buffer[:chunkSize]) + queues[source] = buffer[chunkSize:] } } - if err := writePCM(r.stdin, frame); err != nil { + if err := writePCM(r.stdin, chunk); err != nil { r.setError(err) r.closeEncoder() return @@ -206,19 +252,19 @@ func (r *Recorder) setError(err error) { } } -func NormalizeStereoFrame(samples []int16, frameSize int) []int16 { - out := make([]int16, frameSize*gumble.AudioChannels) - if len(samples) >= frameSize*gumble.AudioChannels && len(samples)%gumble.AudioChannels == 0 { - copy(out, samples[:frameSize*gumble.AudioChannels]) - return out +// NormalizeStereoFrame ensures samples are in stereo interleaved format. +// If stereo is true the samples are returned as-is (already interleaved). +// Mono input is duplicated to both channels. The returned slice preserves +// all input audio without truncation. +func NormalizeStereoFrame(samples []int16, stereo bool) []int16 { + if stereo { + return samples } - limit := frameSize - if len(samples) < limit { - limit = len(samples) - } - for i := 0; i < limit; i++ { - out[i*2] = samples[i] - out[i*2+1] = samples[i] + // Convert mono to stereo by duplicating each sample. + out := make([]int16, len(samples)*gumble.AudioChannels) + for i, s := range samples { + out[i*2] = s + out[i*2+1] = s } return out } @@ -248,7 +294,7 @@ func writePCM(w io.Writer, samples []int16) error { return err } -func ffmpegArgs(format string, path string) []string { +func ffmpegArgs(format string) []string { args := []string{ "-loglevel", "error", "-f", "s16le", @@ -259,5 +305,5 @@ func ffmpegArgs(format string, path string) []string { if format == FormatOpus { args = append(args, "-c:a", "libopus") } - return append(args, "-y", path) + return append(args, "-f", format, "pipe:3") } diff --git a/recording/recorder_test.go b/recording/recorder_test.go index 130d522..10c2fc6 100644 --- a/recording/recorder_test.go +++ b/recording/recorder_test.go @@ -1,12 +1,36 @@ package recording import ( + "io" "os" "path/filepath" + "sync" "testing" "time" ) +type trackingWriteCloser struct{ closed bool } + +func (w *trackingWriteCloser) Write([]byte) (int, error) { return 0, nil } +func (w *trackingWriteCloser) Close() error { w.closed = true; return nil } + +var _ io.WriteCloser = (*trackingWriteCloser)(nil) + +// Regression: Stop closed ffmpeg stdin while the worker could still write, +// creating a spurious closed-pipe recording failure. +func TestStopLeavesEncoderClosureToWorker(t *testing.T) { + stdin := &trackingWriteCloser{} + done := make(chan struct{}) + close(done) + r := &Recorder{stdin: stdin, stop: make(chan struct{}), done: done} + if err := r.Stop(); err != nil { + t.Fatal(err) + } + if stdin.closed { + t.Fatal("Stop closed stdin instead of the worker") + } +} + func TestNormalizeFormat(t *testing.T) { tests := map[string]string{ "": "flac", @@ -35,17 +59,70 @@ func TestUniquePathAvoidsCollision(t *testing.T) { } } +func TestReservePathPreventsConcurrentRecordingCollisions(t *testing.T) { + dir := t.TempDir() + now := time.Date(2026, 5, 14, 12, 30, 0, 0, time.Local) + paths := make(chan string, 2) + errs := make(chan error, 2) + var wg sync.WaitGroup + for range 2 { + wg.Add(1) + go func() { + defer wg.Done() + path, err := reservePath(dir, now, "flac") + if err != nil { + errs <- err + return + } + paths <- path + }() + } + wg.Wait() + close(paths) + close(errs) + for err := range errs { + t.Fatal(err) + } + var reserved []string + for path := range paths { + reserved = append(reserved, path) + } + if len(reserved) != 2 || reserved[0] == reserved[1] { + t.Fatalf("reserved paths = %#v", reserved) + } +} + func TestNormalizeStereoFrame(t *testing.T) { - mono := NormalizeStereoFrame([]int16{1, -2}, 3) - wantMono := []int16{1, 1, -2, -2, 0, 0} + // Mono input duplicating each sample to both channels. + mono := NormalizeStereoFrame([]int16{1, -2, 3}, false) + wantMono := []int16{1, 1, -2, -2, 3, 3} + if len(mono) != len(wantMono) { + t.Fatalf("mono len = %d, want %d", len(mono), len(wantMono)) + } for i := range wantMono { if mono[i] != wantMono[i] { t.Fatalf("mono[%d] = %d, want %d", i, mono[i], wantMono[i]) } } - stereo := NormalizeStereoFrame([]int16{1, 2, 3, 4, 5, 6}, 2) - wantStereo := []int16{1, 2, 3, 4} + // Even-length mono must not be mistaken for stereo. + monoEven := NormalizeStereoFrame([]int16{1, -2}, false) + wantMonoEven := []int16{1, 1, -2, -2} + if len(monoEven) != len(wantMonoEven) { + t.Fatalf("monoEven len = %d, want %d", len(monoEven), len(wantMonoEven)) + } + for i := range wantMonoEven { + if monoEven[i] != wantMonoEven[i] { + t.Fatalf("monoEven[%d] = %d, want %d", i, monoEven[i], wantMonoEven[i]) + } + } + + // Stereo input passes through unchanged. + stereo := NormalizeStereoFrame([]int16{1, 2, 3, 4, 5, 6}, true) + wantStereo := []int16{1, 2, 3, 4, 5, 6} + if len(stereo) != len(wantStereo) { + t.Fatalf("stereo len = %d, want %d", len(stereo), len(wantStereo)) + } for i := range wantStereo { if stereo[i] != wantStereo[i] { t.Fatalf("stereo[%d] = %d, want %d", i, stereo[i], wantStereo[i])