diff --git a/client.go b/client.go index 0b43ac0..639d9df 100644 --- a/client.go +++ b/client.go @@ -57,7 +57,7 @@ func (b *Barnard) connect(reconnect bool) bool { b.Stream.SetAGCEnabled(b.UserConfig.GetAGCEnabled()) // Initialize stereo encoder for file playback - b.Client.AudioEncoderStereo = opus.NewStereoEncoder() + b.Client.SetStereoEncoder(opus.NewStereoEncoder()) // Initialize file player b.FileStreamMutex.Lock() diff --git a/fileplayback/player.go b/fileplayback/player.go index 5f1daf7..4c658b9 100644 --- a/fileplayback/player.go +++ b/fileplayback/player.go @@ -1,6 +1,7 @@ package fileplayback import ( + "context" "encoding/binary" "errors" "io" @@ -10,22 +11,24 @@ import ( "time" "git.stormux.org/storm/barnard/gumble/gumble" - "git.stormux.org/storm/barnard/gumble/go-openal/openal" ) // Player handles file playback and mixing with microphone audio type Player struct { - client *gumble.Client - filename string - audioChan chan gumble.AudioBuffer - stopChan chan struct{} - mutex sync.Mutex - playing bool - errorFunc func(error) + client *gumble.Client + filename string + audioChan chan gumble.AudioBuffer + stopChan chan struct{} + ctx context.Context + cancel context.CancelFunc + cmd *exec.Cmd + mutex sync.Mutex + wg sync.WaitGroup + playing bool + stopping bool + errorFunc func(error) - // Local playback - localSource *openal.Source - localBuffers openal.Buffers + localPlayback func([]byte) } // New creates a new file player @@ -44,6 +47,14 @@ func (p *Player) SetErrorFunc(f func(error)) { p.errorFunc = f } +// SetLocalPlayback sets the callback that plays file audio locally. The +// callback is called with nil when playback stops and should release resources. +func (p *Player) SetLocalPlayback(f func([]byte)) { + p.mutex.Lock() + defer p.mutex.Unlock() + p.localPlayback = f +} + func (p *Player) reportError(err error) { p.mutex.Lock() errorFunc := p.errorFunc @@ -65,17 +76,12 @@ func (p *Player) PlayFile(filename string) error { p.filename = filename - // Initialize local playback - source := openal.NewSource() - p.localSource = &source - p.localSource.SetGain(1.0) - - // Create buffers for local playback - p.localBuffers = openal.NewBuffers(64) - // Start the file reading goroutine p.playing = true + p.stopping = false p.stopChan = make(chan struct{}) + p.ctx, p.cancel = context.WithCancel(context.Background()) + p.wg.Add(1) go p.readFileAudio() return nil @@ -84,31 +90,33 @@ func (p *Player) PlayFile(filename string) error { // Stop stops the currently playing file func (p *Player) Stop() error { p.mutex.Lock() - defer p.mutex.Unlock() - if !p.playing { + p.mutex.Unlock() return errors.New("no file playing") } - - close(p.stopChan) - p.playing = false - - // Clean up local playback - if p.localSource != nil { - p.localSource.Stop() - p.localSource.Delete() - p.localSource = nil - } - if p.localBuffers != nil { - p.localBuffers.Delete() - p.localBuffers = nil + if !p.stopping { + p.stopping = true + close(p.stopChan) + if p.cancel != nil { + p.cancel() + } + terminateProcessGroup(p.cmd) } + p.mutex.Unlock() - // Drain the audio channel + // A new PlayFile must not replace session state until ffmpeg and the old + // worker have exited, otherwise old audio can enter the new playback. + p.wg.Wait() + p.mutex.Lock() + p.playing, p.stopping, p.cancel, p.cmd = false, false, nil, nil + localPlayback := p.localPlayback + p.mutex.Unlock() + if localPlayback != nil { + localPlayback(nil) + } for len(p.audioChan) > 0 { <-p.audioChan } - return nil } @@ -129,37 +137,18 @@ func (p *Player) GetAudioFrame() []int16 { } } -// playLocalAudio plays audio through the local OpenAL source func (p *Player) playLocalAudio(data []byte) { - if p.localSource == nil { - return - } - - // Reclaim processed buffers - if n := p.localSource.BuffersProcessed(); n > 0 { - reclaimedBufs := make(openal.Buffers, n) - p.localSource.UnqueueBuffers(reclaimedBufs) - p.localBuffers = append(p.localBuffers, reclaimedBufs...) - } - - // If we have available buffers, queue more audio - if len(p.localBuffers) > 0 { - buffer := p.localBuffers[len(p.localBuffers)-1] - p.localBuffers = p.localBuffers[:len(p.localBuffers)-1] - - // Set buffer data as stereo - buffer.SetData(openal.FormatStereo16, data, gumble.AudioSampleRate) - p.localSource.QueueBuffer(buffer) - - // Start playing if not already - if p.localSource.State() != openal.Playing { - p.localSource.Play() - } + p.mutex.Lock() + localPlayback := p.localPlayback + p.mutex.Unlock() + if localPlayback != nil { + localPlayback(data) } } // readFileAudio reads audio from the file via ffmpeg func (p *Player) readFileAudio() { + defer p.wg.Done() interval := p.client.Config.AudioInterval frameSize := p.client.Config.AudioFrameSize() @@ -168,7 +157,11 @@ func (p *Player) readFileAudio() { args := []string{"-loglevel", "error", "-i", p.filename} args = append(args, "-ac", "2", "-ar", strconv.Itoa(gumble.AudioSampleRate), "-f", "s16le", "-") - cmd := exec.Command("ffmpeg", args...) + p.mutex.Lock() + ctx := p.ctx + p.mutex.Unlock() + cmd := exec.CommandContext(ctx, "ffmpeg", args...) + configureProcessGroup(cmd) pipe, err := cmd.StdoutPipe() if err != nil { p.mutex.Lock() @@ -185,6 +178,9 @@ func (p *Player) readFileAudio() { p.reportError(errors.New("failed to start ffmpeg: " + err.Error())) return } + p.mutex.Lock() + p.cmd = cmd + p.mutex.Unlock() // Stereo has 2 channels, so we need twice the buffer size byteBuffer := make([]byte, frameSize*2*2) // frameSize * 2 channels * 2 bytes per sample @@ -195,28 +191,27 @@ func (p *Player) readFileAudio() { for { select { case <-p.stopChan: - cmd.Process.Kill() + terminateProcessGroup(cmd) cmd.Wait() return case <-ticker.C: n, err := io.ReadFull(pipe, byteBuffer) if err != nil || n != len(byteBuffer) { - // File finished playing + select { + case <-p.stopChan: + cmd.Wait() + return + default: + } + // File finished playing. p.mutex.Lock() p.playing = false - // Clean up local playback - if p.localSource != nil { - p.localSource.Stop() - p.localSource.Delete() - p.localSource = nil - } - if p.localBuffers != nil { - p.localBuffers.Delete() - p.localBuffers = nil - } + localPlayback := p.localPlayback p.mutex.Unlock() + if localPlayback != nil { + localPlayback(nil) + } cmd.Wait() - // Notify that file finished p.reportError(errors.New("file playback finished")) return } diff --git a/fileplayback/player_regression_test.go b/fileplayback/player_regression_test.go new file mode 100644 index 0000000..347a113 --- /dev/null +++ b/fileplayback/player_regression_test.go @@ -0,0 +1,24 @@ +package fileplayback + +import ( + "testing" + "time" +) + +// Regression: Stop returned before the previous ffmpeg worker exited, allowing +// a subsequent PlayFile to replace shared state while old audio was still sent. +func TestStopWaitsForPlaybackWorker(t *testing.T) { + p := &Player{playing: true, stopChan: make(chan struct{})} + p.wg.Add(1) + done := make(chan error, 1) + go func() { done <- p.Stop() }() + select { + case <-done: + t.Fatal("Stop returned before playback worker ended") + case <-time.After(20 * time.Millisecond): + } + p.wg.Done() + if err := <-done; err != nil { + t.Fatal(err) + } +} diff --git a/fileplayback/process_other.go b/fileplayback/process_other.go new file mode 100644 index 0000000..744f15e --- /dev/null +++ b/fileplayback/process_other.go @@ -0,0 +1,13 @@ +//go:build !(aix || darwin || dragonfly || freebsd || linux || netbsd || openbsd || solaris) + +package fileplayback + +import "os/exec" + +func configureProcessGroup(cmd *exec.Cmd) {} + +func terminateProcessGroup(cmd *exec.Cmd) { + if cmd != nil && cmd.Process != nil { + _ = cmd.Process.Kill() + } +} diff --git a/fileplayback/process_unix.go b/fileplayback/process_unix.go new file mode 100644 index 0000000..00b15bf --- /dev/null +++ b/fileplayback/process_unix.go @@ -0,0 +1,19 @@ +//go:build aix || darwin || dragonfly || freebsd || linux || netbsd || openbsd || solaris + +package fileplayback + +import ( + "os/exec" + "syscall" +) + +func configureProcessGroup(cmd *exec.Cmd) { + cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} +} + +func terminateProcessGroup(cmd *exec.Cmd) { + if cmd == nil || cmd.Process == nil { + return + } + _ = syscall.Kill(-cmd.Process.Pid, syscall.SIGKILL) +} diff --git a/fileplayback/process_unix_test.go b/fileplayback/process_unix_test.go new file mode 100644 index 0000000..3fed92a --- /dev/null +++ b/fileplayback/process_unix_test.go @@ -0,0 +1,16 @@ +//go:build aix || darwin || dragonfly || freebsd || linux || netbsd || openbsd || solaris + +package fileplayback + +import ( + "os/exec" + "testing" +) + +func TestConfigureProcessGroupCreatesSeparateGroup(t *testing.T) { + cmd := exec.Command("true") + configureProcessGroup(cmd) + if cmd.SysProcAttr == nil || !cmd.SysProcAttr.Setpgid { + t.Fatal("ffmpeg process was not configured to lead its own process group") + } +} diff --git a/gumble/gumble/client.go b/gumble/gumble/client.go index 7b05c1a..f7585f9 100644 --- a/gumble/gumble/client.go +++ b/gumble/gumble/client.go @@ -396,6 +396,13 @@ func (c *Client) Send(message Message) { message.writeMessage(c) } +// SetStereoEncoder installs the encoder used for stereo file playback. +func (c *Client) SetStereoEncoder(encoder AudioEncoder) { + c.volatile.Lock() + defer c.volatile.Unlock() + c.AudioEncoderStereo = encoder +} + // EnableStereoEncoder switches to stereo encoding for file playback. func (c *Client) EnableStereoEncoder() { c.volatile.Lock() @@ -448,11 +455,16 @@ func (c *Client) UDPActive() bool { return c.udpActive } -// DisableStereoEncoder switches back to mono encoding for voice. +// DisableStereoEncoder switches back to mono encoding for voice and +// resets the stereo encoder so stale state does not bleed into the +// next file playback. func (c *Client) DisableStereoEncoder() { c.volatile.Lock() defer c.volatile.Unlock() c.useStereoEncoder = false + if c.AudioEncoderStereo != nil { + c.AudioEncoderStereo.Reset() + } } // IsStereoEncoderEnabled returns true if stereo encoding is currently active. 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") + } +}