diff --git a/gumble/gumble/audio.go b/gumble/gumble/audio.go index 74e868d..b99ad67 100644 --- a/gumble/gumble/audio.go +++ b/gumble/gumble/audio.go @@ -92,6 +92,10 @@ type AudioPacket struct { AudioBuffer + // Terminator marks the final packet in a talk burst. Audio listeners use + // it to discard ordering state before the sender starts a new burst. + Terminator bool + HasPosition bool X, Y, Z float32 VolumeAdjustment float32 diff --git a/gumble/gumble/handlers.go b/gumble/gumble/handlers.go index e920896..ed31bf2 100644 --- a/gumble/gumble/handlers.go +++ b/gumble/gumble/handlers.go @@ -210,6 +210,7 @@ func (c *Client) handleUDPTunnel(buffer []byte) error { }, Sequence: seq, AudioBuffer: AudioBuffer(pcm), + Terminator: isFinal, } if len(buffer)-audioLength == 3*4 { @@ -223,6 +224,10 @@ func (c *Client) handleUDPTunnel(buffer []byte) error { } c.dispatchAudio(user, &event) + if isFinal { + decoder.Reset() + user.audioSequenceValid = false + } return nil } diff --git a/gumble/gumble/udp15.go b/gumble/gumble/udp15.go index 2b742dc..83977ae 100644 --- a/gumble/gumble/udp15.go +++ b/gumble/gumble/udp15.go @@ -752,6 +752,9 @@ func (c *Client) dispatchOpus15(pktNum uint64, session uint32, frameNum int64, o decoder.Reset() user.audioSequenceValid = false user.audioFrameStep = 0 + // The audio stream remains open between talk bursts. Deliver the + // terminator so listeners can reset their own packet ordering state. + c.dispatchAudio(user, &AudioPacket{Client: c, Sender: user, Terminator: true}) log.Info("UDP15 #%d: terminator for %s, decoder reset", pktNum, user.Name) return } @@ -760,7 +763,7 @@ func (c *Client) dispatchOpus15(pktNum uint64, session uint32, frameNum int64, o return } - c.decodeAndDispatch(pktNum, user, decoder, frameNum, opusData, context, position, volumeAdjustment) + c.decodeAndDispatch(pktNum, user, decoder, frameNum, opusData, terminator, context, position, volumeAdjustment) } // handleLegacyUDPVoice parses the legacy UDPVoice format (type byte 0x80) @@ -809,7 +812,7 @@ func (c *Client) handleLegacyUDPVoice(pktNum uint64, data []byte) { } // decodeAndDispatch decodes an Opus frame and dispatches PCM to audio listeners. -func (c *Client) decodeAndDispatch(pktNum uint64, user *User, decoder AudioDecoder, frameNum int64, opusData []byte, context uint32, position *[3]float32, volumeAdjustment float32) { +func (c *Client) decodeAndDispatch(pktNum uint64, user *User, decoder AudioDecoder, frameNum int64, opusData []byte, terminator bool, context uint32, position *[3]float32, volumeAdjustment float32) { // Frame numbers are timestamps in 10 ms units, not packet counters. For // example, a standard 20 ms Opus packet advances its frame number by two. // Only generate PLC for complete missing packets; treating every timestamp @@ -858,6 +861,7 @@ func (c *Client) decodeAndDispatch(pktNum uint64, user *User, decoder AudioDecod Target: &VoiceTarget{ID: context}, Sequence: frameNum, AudioBuffer: AudioBuffer(pcm), + Terminator: terminator, VolumeAdjustment: volumeAdjustment, } if position != nil { @@ -865,6 +869,11 @@ func (c *Client) decodeAndDispatch(pktNum uint64, user *User, decoder AudioDecod event.X, event.Y, event.Z = position[0], position[1], position[2] } c.dispatchAudio(user, &event) + if terminator { + decoder.Reset() + user.audioSequenceValid = false + user.audioFrameStep = 0 + } } // missingAudioPackets returns the number of whole packets absent from a diff --git a/gumble/gumble/udp15_terminator_test.go b/gumble/gumble/udp15_terminator_test.go new file mode 100644 index 0000000..501ba80 --- /dev/null +++ b/gumble/gumble/udp15_terminator_test.go @@ -0,0 +1,37 @@ +package gumble + +import "testing" + +type terminatorDecoder struct{ resets int } + +func (d *terminatorDecoder) ID() int { return audioCodecIDOpus } +func (d *terminatorDecoder) Decode([]byte, int) ([]int16, error) { return nil, nil } +func (d *terminatorDecoder) Reset() { d.resets++ } + +type terminatorListener struct{ packets chan *AudioPacket } + +func (l *terminatorListener) OnAudioStream(e *AudioStreamEvent) { + go func() { l.packets <- <-e.C }() +} + +func TestUDP15EmptyTerminatorResetsAudioListeners(t *testing.T) { + decoder := &terminatorDecoder{} + listener := &terminatorListener{packets: make(chan *AudioPacket, 1)} + config := NewConfig() + config.AttachAudio(listener) + user := &User{Session: 1, Name: "speaker", decoder: decoder, audioSequenceValid: true} + client := &Client{Config: config, Users: Users{user.Session: user}} + + client.dispatchOpus15(1, user.Session, 0, nil, true, 0, nil, 0) + + packet := <-listener.packets + if !packet.Terminator { + t.Fatal("empty UDP terminator was not delivered to audio listeners") + } + if packet.AudioBuffer != nil { + t.Fatalf("terminator carried unexpected audio: %v", packet.AudioBuffer) + } + if decoder.resets != 1 || user.audioSequenceValid { + t.Fatalf("terminator did not reset decoder state: resets=%d valid=%v", decoder.resets, user.audioSequenceValid) + } +} diff --git a/gumble/gumbleopenal/stream.go b/gumble/gumbleopenal/stream.go index bb691fa..91a6a8f 100644 --- a/gumble/gumbleopenal/stream.go +++ b/gumble/gumbleopenal/stream.go @@ -487,6 +487,12 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) { var jitterNextSeq int64 var jitterInit, jitterStarted bool var jitterDrainLogCounter, jitterAnomalyLogCounter int + resetJitter := func() { + jitterBuf = nil + jitterNextSeq = 0 + jitterInit = false + jitterStarted = false + } // insertSorted inserts a packet into the jitter buffer sorted // by sequence number. @@ -541,6 +547,14 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) { } for packet := range e.C { + // A talk burst may restart its frame numbers from zero. Reset before + // testing local mute so an unmute cannot retain the previous burst's + // timestamp and discard the new burst as permanently late. + if packet.Terminator { + resetJitter() + continue + } + // Skip processing if user is locally muted if e.User.LocallyMuted() { continue