diff --git a/fix.txt b/fix.txt index 2a8cde9..ffb60c6 100644 --- a/fix.txt +++ b/fix.txt @@ -121,7 +121,7 @@ Priority 1: transport, lifecycle, and correctness output file creation fails, the tone goroutine continues. Create the saver first, or close/wait for the tone generator and reset Tx on every failure. -15. Gumble ffmpeg Pause can block forever +[x] 15. Gumble ffmpeg Pause can block forever File: gumble/gumbleffmpeg/stream.go Pause checks StatePlaying, releases the lock, then sends on an unbuffered pause channel. If process exits in between, no receiver remains. Redesign diff --git a/gumble/gumbleffmpeg/stream.go b/gumble/gumbleffmpeg/stream.go index e5518fe..bed0c14 100644 --- a/gumble/gumbleffmpeg/stream.go +++ b/gumble/gumbleffmpeg/stream.go @@ -57,7 +57,7 @@ func New(client *gumble.Client, source Source) *Stream { Volume: 1.0, Source: source, Command: "ffmpeg", - pause: make(chan struct{}), + pause: make(chan struct{}, 1), state: StateInitial, } } @@ -124,7 +124,12 @@ func (s *Stream) Pause() error { } s.state = StatePaused s.l.Unlock() - s.pause <- struct{}{} + // The process can exit after the state check. A buffered, coalesced pause + // request preserves the state transition without blocking the caller. + select { + case s.pause <- struct{}{}: + default: + } return nil } @@ -176,9 +181,12 @@ func (s *Stream) process() { return } int16Buffer := make([]int16, frameSize) + s.l.Lock() + volume := s.Volume + s.l.Unlock() for i := range int16Buffer { float := float32(int16(binary.LittleEndian.Uint16(byteBuffer[i*2 : (i+1)*2]))) - int16Buffer[i] = int16(s.Volume * float) + int16Buffer[i] = int16(volume * float) } atomic.AddInt64(&s.elapsed, int64(interval)) outgoing <- gumble.AudioBuffer(int16Buffer) diff --git a/gumble/gumbleffmpeg/stream_test.go b/gumble/gumbleffmpeg/stream_test.go new file mode 100644 index 0000000..cd1300c --- /dev/null +++ b/gumble/gumbleffmpeg/stream_test.go @@ -0,0 +1,22 @@ +package gumbleffmpeg + +import ( + "testing" + "time" +) + +// Regression: Pause sent on an unbuffered channel after the process had +// exited, leaving callers blocked forever. +func TestPauseDoesNotBlockWhenProcessHasExited(t *testing.T) { + s := &Stream{state: StatePlaying, pause: make(chan struct{}, 1)} + done := make(chan error, 1) + go func() { done <- s.Pause() }() + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(time.Second): + t.Fatal("Pause blocked after process exit") + } +}