Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8a8cfde01d | ||
|
|
c0ef036934 | ||
|
|
d338925da6 | ||
|
|
893908f9a1 | ||
|
|
05dc6e4e0e | ||
|
|
883f7250f5 | ||
|
|
4510c25350 | ||
|
|
d4f8d56c8e | ||
|
|
c77d4bac3e | ||
|
|
0f34831c5b | ||
|
|
4c5a54c2dd | ||
|
|
eae1d8b99a |
@@ -140,6 +140,35 @@ If you modify the config file while Barnard is running, your changes may be over
|
||||
You can set username and defaultserver in your config file, and they will be used if none is specified when launching barnard.
|
||||
(Note that the default username (an empty string) and the default server name (localhost:64738) have been the defaults for barnard up to this point, and have been left that way for compatibility.)
|
||||
|
||||
## Audio Packet Duration
|
||||
|
||||
Barnard sends 10 ms audio packets by default. On a slow or unstable connection,
|
||||
using larger packets can reduce packet overhead and make short dropouts less
|
||||
noticeable, at the cost of additional voice latency. Start Barnard with one of
|
||||
the supported durations:
|
||||
|
||||
```sh
|
||||
barnard --audio-interval 20
|
||||
```
|
||||
|
||||
Supported values are `10`, `20`, `40`, and `60` milliseconds. Try `20` ms
|
||||
first; use `40` ms only if the connection remains unreliable.
|
||||
|
||||
## Incoming Audio Jitter Buffer
|
||||
|
||||
Barnard holds 40 ms of audio separately for each speaker before starting
|
||||
playback. This prevents brief delayed UDP packets from draining OpenAL's audio
|
||||
queue, which otherwise produces clicks or pops. To adjust this tradeoff between
|
||||
resilience and added incoming latency:
|
||||
|
||||
```sh
|
||||
barnard --jitter-buffer 60
|
||||
```
|
||||
|
||||
Supported values are `0`, `20`, `40` (default), and `60` milliseconds. Try
|
||||
`60` ms for a lossy or jittery connection. Use `0` only when minimizing latency
|
||||
is more important than avoiding playback underruns.
|
||||
|
||||
## Audio Devices
|
||||
|
||||
You can set the default input and output devices in the config file as well.
|
||||
|
||||
@@ -70,6 +70,38 @@ func TestNotificationExpansionIsSinglePassAndNotifyDoesNotBlock(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAudioIntervalDuration(t *testing.T) {
|
||||
for _, milliseconds := range []int{10, 20, 40, 60} {
|
||||
got, err := audioIntervalDuration(milliseconds)
|
||||
if err != nil {
|
||||
t.Errorf("audioIntervalDuration(%d): %v", milliseconds, err)
|
||||
continue
|
||||
}
|
||||
if got != time.Duration(milliseconds)*time.Millisecond {
|
||||
t.Errorf("audioIntervalDuration(%d) = %v", milliseconds, got)
|
||||
}
|
||||
}
|
||||
if _, err := audioIntervalDuration(30); err == nil {
|
||||
t.Fatal("audioIntervalDuration accepted unsupported duration")
|
||||
}
|
||||
}
|
||||
|
||||
func TestJitterBufferDuration(t *testing.T) {
|
||||
for _, milliseconds := range []int{0, 20, 40, 60} {
|
||||
got, err := jitterBufferDuration(milliseconds)
|
||||
if err != nil {
|
||||
t.Errorf("jitterBufferDuration(%d): %v", milliseconds, err)
|
||||
continue
|
||||
}
|
||||
if got != time.Duration(milliseconds)*time.Millisecond {
|
||||
t.Errorf("jitterBufferDuration(%d) = %v", milliseconds, got)
|
||||
}
|
||||
}
|
||||
if _, err := jitterBufferDuration(10); err == nil {
|
||||
t.Fatal("jitterBufferDuration accepted unsupported duration")
|
||||
}
|
||||
}
|
||||
|
||||
func TestServerAddressDefaultsPortWithoutBreakingIPv6(t *testing.T) {
|
||||
for input, want := range map[string]string{
|
||||
"server": "server:64738",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -3,6 +3,7 @@ package gumble
|
||||
import (
|
||||
"crypto/tls"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math"
|
||||
"net"
|
||||
"runtime"
|
||||
@@ -101,6 +102,19 @@ func Dial(config *Config) (*Client, error) {
|
||||
return DialWithDialer(new(net.Dialer), config, nil)
|
||||
}
|
||||
|
||||
// tlsServerName returns the hostname portion of a Mumble server address for
|
||||
// TLS certificate verification and SNI.
|
||||
func tlsServerName(address string) (string, error) {
|
||||
host, _, err := net.SplitHostPort(address)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("gumble: derive TLS server name from %q: %w", address, err)
|
||||
}
|
||||
if host == "" {
|
||||
return "", fmt.Errorf("gumble: derive TLS server name from %q: empty host", address)
|
||||
}
|
||||
return host, nil
|
||||
}
|
||||
|
||||
// DialWithDialer connects to the Mumble server at the address given in config.
|
||||
//
|
||||
// The function returns after the connection has been established, the initial
|
||||
@@ -120,6 +134,23 @@ func DialWithDialer(dialer *net.Dialer, config *Config, tlsConfig *tls.Config) (
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// tls.Client cannot infer a server name from an already-open connection.
|
||||
// Clone the caller's configuration before deriving it so reconnects and
|
||||
// concurrent clients do not mutate a shared configuration.
|
||||
if tlsConfig == nil {
|
||||
tlsConfig = &tls.Config{}
|
||||
} else {
|
||||
tlsConfig = tlsConfig.Clone()
|
||||
}
|
||||
if tlsConfig.ServerName == "" {
|
||||
serverName, err := tlsServerName(config.Address)
|
||||
if err != nil {
|
||||
rawConn.Close()
|
||||
return nil, err
|
||||
}
|
||||
tlsConfig.ServerName = serverName
|
||||
}
|
||||
conn := tls.Client(rawConn, tlsConfig)
|
||||
// net.Dialer.Timeout covers only the TCP dial. Apply the same bounded
|
||||
// deadline to TLS negotiation so a peer that accepts but never responds
|
||||
|
||||
@@ -0,0 +1,28 @@
|
||||
package gumble
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestTLSServerNameUsesAddressHost(t *testing.T) {
|
||||
for _, test := range []struct {
|
||||
address string
|
||||
want string
|
||||
}{
|
||||
{"mumble.example:64738", "mumble.example"},
|
||||
{"[2001:db8::1]:64738", "2001:db8::1"},
|
||||
} {
|
||||
got, err := tlsServerName(test.address)
|
||||
if err != nil {
|
||||
t.Errorf("tlsServerName(%q): %v", test.address, err)
|
||||
continue
|
||||
}
|
||||
if got != test.want {
|
||||
t.Errorf("tlsServerName(%q) = %q, want %q", test.address, got, test.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestTLSServerNameRejectsAddressWithoutHost(t *testing.T) {
|
||||
if _, err := tlsServerName(":64738"); err == nil {
|
||||
t.Fatal("tlsServerName accepted an empty host")
|
||||
}
|
||||
}
|
||||
+10
-3
@@ -25,6 +25,9 @@ type Config struct {
|
||||
AudioInterval time.Duration
|
||||
// AudioDataBytes is the number of bytes that an audio frame can use.
|
||||
AudioDataBytes int
|
||||
// IncomingAudioBuffer is the amount of per-speaker audio retained before
|
||||
// playback starts, absorbing jitter in incoming UDP packet delivery.
|
||||
IncomingAudioBuffer time.Duration
|
||||
|
||||
// DisableUDP forces all audio to use the TCP tunnel instead of UDP.
|
||||
DisableUDP bool
|
||||
@@ -38,9 +41,10 @@ type Config struct {
|
||||
// NewConfig returns a new Config struct with default values set.
|
||||
func NewConfig() *Config {
|
||||
return &Config{
|
||||
Buffers: 8,
|
||||
AudioInterval: AudioDefaultInterval,
|
||||
AudioDataBytes: AudioDefaultDataBytes,
|
||||
Buffers: 8,
|
||||
AudioInterval: AudioDefaultInterval,
|
||||
AudioDataBytes: AudioDefaultDataBytes,
|
||||
IncomingAudioBuffer: 40 * time.Millisecond,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -54,6 +58,9 @@ func (c *Config) Validate() error {
|
||||
if c.AudioDataBytes <= 0 {
|
||||
return fmt.Errorf("gumble: AudioDataBytes must be positive")
|
||||
}
|
||||
if c.IncomingAudioBuffer < 0 {
|
||||
return fmt.Errorf("gumble: IncomingAudioBuffer must not be negative")
|
||||
}
|
||||
if c.Buffers <= 0 {
|
||||
return fmt.Errorf("gumble: Buffers must be positive")
|
||||
}
|
||||
|
||||
@@ -223,6 +223,11 @@ func (c *Client) handleUDPTunnel(buffer []byte) error {
|
||||
}
|
||||
|
||||
c.dispatchAudio(user, &event)
|
||||
if isFinal {
|
||||
decoder.Reset()
|
||||
user.audioSequenceValid = false
|
||||
c.dispatchAudio(user, &AudioPacket{Client: c, Sender: user, Terminator: true})
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -59,7 +59,10 @@ func (c *Client) udpReadRoutine() {
|
||||
return
|
||||
}
|
||||
packetCount++
|
||||
if log.Enabled(log.LevelDebug) {
|
||||
// A synchronous log write for every UDP datagram can itself make the
|
||||
// reader fall behind and lose voice packets. Keep enough samples to
|
||||
// diagnose framing while avoiding work on the audio hot path.
|
||||
if log.Enabled(log.LevelDebug) && (packetCount <= 3 || packetCount%1000 == 0) {
|
||||
log.Debug("UDP recv #%d: %d bytes from %s hex=%s",
|
||||
packetCount, n, addr, hex.EncodeToString(buf[:n]))
|
||||
}
|
||||
|
||||
+37
-9
@@ -292,10 +292,13 @@ func (cs *cryptState15) decrypt15(packet []byte) ([]byte, error) {
|
||||
backupIV(cs.decryptIV[:])
|
||||
restore = true
|
||||
} else if ivByte > cs.decryptIV[0] && diff > 0 {
|
||||
// We missed packets; catch up. Already handled above.
|
||||
// We missed packets; move the low IV byte forward.
|
||||
cs.decryptIV[0] = ivByte
|
||||
} else if ivByte < cs.decryptIV[0] && diff > 0 {
|
||||
// Wrapped forward; advance and catch up.
|
||||
advanceIV(cs.decryptIV[:])
|
||||
// We missed packets across a low-byte wrap. The IV's higher
|
||||
// bytes must advance even though the received low byte is set
|
||||
// below rather than incremented.
|
||||
advanceIVHighBytes(cs.decryptIV[:])
|
||||
cs.decryptIV[0] = ivByte
|
||||
} else {
|
||||
return nil, errors.New("gumble: OCB IV too far off")
|
||||
@@ -340,6 +343,16 @@ func advanceIV(iv []byte) {
|
||||
}
|
||||
}
|
||||
|
||||
// advanceIVHighBytes advances all but the low IV byte as a little-endian integer.
|
||||
func advanceIVHighBytes(iv []byte) {
|
||||
for i := 1; i < len(iv); i++ {
|
||||
iv[i]++
|
||||
if iv[i] != 0 {
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// backupIV decrements a 16-byte IV as a little-endian integer.
|
||||
func backupIV(iv []byte) {
|
||||
for i := 0; i < len(iv); i++ {
|
||||
@@ -626,8 +639,10 @@ func (c *Client) WriteAudioUDP15(format byte, target uint32, sequence int64, dat
|
||||
return false, err
|
||||
}
|
||||
|
||||
log.Debug("UDP15 send: frame=%d opus_len=%d enc_len=%d final=%v",
|
||||
frameNum, len(data), len(encrypted), final)
|
||||
if log.Enabled(log.LevelDebug) && (frameNum < 3 || frameNum%1000 == 0 || final) {
|
||||
log.Debug("UDP15 send: frame=%d opus_len=%d enc_len=%d final=%v",
|
||||
frameNum, len(data), len(encrypted), final)
|
||||
}
|
||||
|
||||
_, err = udpConn.Write(encrypted)
|
||||
if err != nil {
|
||||
@@ -662,7 +677,9 @@ func (c *Client) HandleUDPPacket15(packet []byte, pktNum uint64) {
|
||||
return
|
||||
}
|
||||
|
||||
log.Debug("UDP15 #%d: decrypt OK, plaintext_len=%d", pktNum, len(plaintext))
|
||||
if log.Enabled(log.LevelDebug) && (pktNum <= 3 || pktNum%1000 == 0) {
|
||||
log.Debug("UDP15 #%d: decrypt OK, plaintext_len=%d", pktNum, len(plaintext))
|
||||
}
|
||||
c.markUDPActive()
|
||||
|
||||
// Check type byte (0x00 = Audio, 0x01 = Ping)
|
||||
@@ -735,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
|
||||
}
|
||||
@@ -743,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)
|
||||
@@ -792,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
|
||||
@@ -828,7 +848,9 @@ func (c *Client) decodeAndDispatch(pktNum uint64, user *User, decoder AudioDecod
|
||||
return
|
||||
}
|
||||
|
||||
log.Debug("UDP15 #%d: Opus OK for %s, pcm_samples=%d", pktNum, user.Name, len(pcm))
|
||||
if log.Enabled(log.LevelDebug) && (pktNum <= 3 || pktNum%1000 == 0) {
|
||||
log.Debug("UDP15 #%d: Opus OK for %s, pcm_samples=%d", pktNum, user.Name, len(pcm))
|
||||
}
|
||||
user.audioSequence = frameNum
|
||||
user.audioSequenceValid = true
|
||||
user.audioFrameStep = audioFrameStep(len(pcm))
|
||||
@@ -846,6 +868,12 @@ 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
|
||||
c.dispatchAudio(user, &AudioPacket{Client: c, Sender: user, Terminator: true})
|
||||
}
|
||||
}
|
||||
|
||||
// missingAudioPackets returns the number of whole packets absent from a
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -32,6 +32,75 @@ func TestAdvanceIV(t *testing.T) {
|
||||
}
|
||||
|
||||
// Regression coverage for the native IV carry path at the 255->256 wrap.
|
||||
func TestCryptState15DecryptsAfterMissedPackets(t *testing.T) {
|
||||
key := mustDecodeHex("93360b0f86a926c4561563469026eb94")
|
||||
nonce := mustDecodeHex("10000000000000000000000000000000")
|
||||
out, in := &cryptState15{}, &cryptState15{}
|
||||
if err := out.setup15(key, nonce, nonce); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := in.setup15(key, nonce, nonce); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
first, err := out.encrypt15([]byte("first"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := in.decrypt15(first); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i := 0; i < 3; i++ {
|
||||
if _, err := out.encrypt15([]byte("dropped")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
last, err := out.encrypt15([]byte("after loss"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
plain, err := in.decrypt15(last)
|
||||
if err != nil || !bytes.Equal(plain, []byte("after loss")) {
|
||||
t.Fatalf("decrypt after missed packets = %q, %v", plain, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCryptState15DecryptsAfterMissedPacketsAcrossIVByteWrap(t *testing.T) {
|
||||
key := mustDecodeHex("93360b0f86a926c4561563469026eb94")
|
||||
nonce := mustDecodeHex("fa000000000000000000000000000000")
|
||||
out, in := &cryptState15{}, &cryptState15{}
|
||||
if err := out.setup15(key, nonce, nonce); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := in.setup15(key, nonce, nonce); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
first, err := out.encrypt15([]byte("before wrap"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := in.decrypt15(first); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i := 0; i < 6; i++ {
|
||||
if _, err := out.encrypt15([]byte("dropped")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
last, err := out.encrypt15([]byte("after wrap"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
plain, err := in.decrypt15(last)
|
||||
if err != nil || !bytes.Equal(plain, []byte("after wrap")) {
|
||||
t.Fatalf("decrypt after missed packets across IV wrap = %q, %v", plain, err)
|
||||
}
|
||||
if in.decryptIV[0] != 2 || in.decryptIV[1] != 1 {
|
||||
t.Fatalf("unexpected IV after wrapped loss: %x", in.decryptIV[:2])
|
||||
}
|
||||
}
|
||||
|
||||
func TestCryptState15DecryptsAcrossIVByteWrap(t *testing.T) {
|
||||
key := mustDecodeHex("93360b0f86a926c4561563469026eb94")
|
||||
clientNonce := mustDecodeHex("ff000000000000000000000000000000")
|
||||
|
||||
@@ -52,14 +52,22 @@ const recorderOutgoingSource uint32 = ^uint32(0)
|
||||
|
||||
const (
|
||||
maxBufferSize = 11520 // Max frame size (2880) * bytes per stereo sample (4)
|
||||
jitterMinPackets = 3
|
||||
jitterMaxPackets = 10
|
||||
jitterMaxPackets = 50
|
||||
)
|
||||
|
||||
// jitterPlaybackReady holds the initial playout delay only once. Requiring
|
||||
// the minimum on every packet drains and refills the renderer in bursts.
|
||||
func jitterPlaybackReady(started bool, buffered int) bool {
|
||||
return started || buffered >= jitterMinPackets
|
||||
// jitterPlaybackReady holds the requested initial playout delay only once.
|
||||
// Requiring the delay on every packet drains and refills the renderer in bursts.
|
||||
func jitterPlaybackReady(started bool, buffered, target time.Duration) bool {
|
||||
return started || buffered >= target
|
||||
}
|
||||
|
||||
func audioPacketDuration(packet *gumble.AudioPacket) time.Duration {
|
||||
if packet == nil || len(packet.AudioBuffer) == 0 {
|
||||
return 0
|
||||
}
|
||||
// Opus decoders deliver interleaved stereo PCM to this renderer.
|
||||
frames := len(packet.AudioBuffer) / gumble.AudioChannels
|
||||
return time.Duration(frames) * time.Second / gumble.AudioSampleRate
|
||||
}
|
||||
|
||||
var (
|
||||
@@ -484,8 +492,17 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
// Jitter buffer: collects incoming packets, reorders by
|
||||
// sequence number, and releases them after a small initial delay.
|
||||
var jitterBuf []*gumble.AudioPacket
|
||||
var jitterDuration time.Duration
|
||||
var jitterNextSeq int64
|
||||
var jitterInit, jitterStarted bool
|
||||
var jitterDrainLogCounter, jitterAnomalyLogCounter int
|
||||
resetJitter := func() {
|
||||
jitterBuf = nil
|
||||
jitterDuration = 0
|
||||
jitterNextSeq = 0
|
||||
jitterInit = false
|
||||
jitterStarted = false
|
||||
}
|
||||
|
||||
// insertSorted inserts a packet into the jitter buffer sorted
|
||||
// by sequence number.
|
||||
@@ -506,6 +523,7 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
jitterBuf = append(jitterBuf, nil)
|
||||
copy(jitterBuf[i+1:], jitterBuf[i:])
|
||||
jitterBuf[i] = p
|
||||
jitterDuration += audioPacketDuration(p)
|
||||
}
|
||||
|
||||
// popNext removes and returns the packet with the expected next
|
||||
@@ -516,6 +534,7 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
}
|
||||
p := jitterBuf[0]
|
||||
jitterBuf = jitterBuf[1:]
|
||||
jitterDuration -= audioPacketDuration(p)
|
||||
// Frame numbers are Mumble timestamps in 10 ms units.
|
||||
// Compute the actual step from the PCM sample count so we
|
||||
// never skip a legitimate gap.
|
||||
@@ -540,6 +559,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
|
||||
@@ -556,14 +583,13 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
|
||||
// Hold only the initial packets. Once playback starts, drain every
|
||||
// ready packet so the renderer is fed continuously rather than in
|
||||
// bursts of jitterMinPackets packets.
|
||||
if !jitterPlaybackReady(jitterStarted, len(jitterBuf)) {
|
||||
// bursts of packets.
|
||||
if !jitterPlaybackReady(jitterStarted, jitterDuration, e.Client.Config.IncomingAudioBuffer) {
|
||||
continue
|
||||
}
|
||||
jitterStarted = true
|
||||
|
||||
// Drain all packets that are ready (in sequence order)
|
||||
drainedCount := 0
|
||||
for {
|
||||
pkt := popNext()
|
||||
if pkt == nil {
|
||||
@@ -571,16 +597,23 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
if jitterBuf[0].Sequence < jitterNextSeq {
|
||||
// Late or duplicate: discard so it doesn't
|
||||
// permanently block the drain loop.
|
||||
log.Debug("jitter: discarding late seq=%d for %s (next=%d buf=%d)",
|
||||
jitterBuf[0].Sequence, e.User.Name, jitterNextSeq, len(jitterBuf))
|
||||
jitterAnomalyLogCounter++
|
||||
if jitterAnomalyLogCounter <= 3 || jitterAnomalyLogCounter%1000 == 0 {
|
||||
log.Debug("jitter: discarding late seq=%d for %s (next=%d buf=%d)",
|
||||
jitterBuf[0].Sequence, e.User.Name, jitterNextSeq, len(jitterBuf))
|
||||
}
|
||||
jitterDuration -= audioPacketDuration(jitterBuf[0])
|
||||
jitterBuf = jitterBuf[1:]
|
||||
continue
|
||||
}
|
||||
if jitterBuf[0].Sequence > jitterNextSeq {
|
||||
// Gap in sequence: skip ahead so we don't
|
||||
// wait forever for a lost packet.
|
||||
log.Debug("jitter: seq gap for %s, skipping from %d to %d (buf=%d)",
|
||||
e.User.Name, jitterNextSeq, jitterBuf[0].Sequence, len(jitterBuf))
|
||||
jitterAnomalyLogCounter++
|
||||
if jitterAnomalyLogCounter <= 3 || jitterAnomalyLogCounter%1000 == 0 {
|
||||
log.Debug("jitter: seq gap for %s, skipping from %d to %d (buf=%d)",
|
||||
e.User.Name, jitterNextSeq, jitterBuf[0].Sequence, len(jitterBuf))
|
||||
}
|
||||
jitterNextSeq = jitterBuf[0].Sequence
|
||||
continue
|
||||
}
|
||||
@@ -589,8 +622,8 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
}
|
||||
break
|
||||
}
|
||||
drainedCount++
|
||||
if drainedCount <= 3 || drainedCount%50 == 0 {
|
||||
jitterDrainLogCounter++
|
||||
if jitterDrainLogCounter <= 3 || jitterDrainLogCounter%1000 == 0 {
|
||||
log.Debug("jitter: draining seq=%d for %s (buf=%d emptyBufs=%d)",
|
||||
pkt.Sequence, e.User.Name, len(jitterBuf), len(emptyBufs))
|
||||
}
|
||||
|
||||
@@ -4,8 +4,10 @@ import (
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.stormux.org/storm/barnard/gumble/go-openal/openal"
|
||||
"git.stormux.org/storm/barnard/gumble/gumble"
|
||||
)
|
||||
|
||||
// Regression: audio cleanup could send a final render command after Destroy
|
||||
@@ -49,17 +51,24 @@ func TestStopSourceWaitsForWorker(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestJitterPlaybackDelayAppliesOnlyAtStartup(t *testing.T) {
|
||||
if jitterPlaybackReady(false, jitterMinPackets-1) {
|
||||
if jitterPlaybackReady(false, 20*time.Millisecond, 40*time.Millisecond) {
|
||||
t.Fatal("jitter playback started before initial buffer filled")
|
||||
}
|
||||
if !jitterPlaybackReady(false, jitterMinPackets) {
|
||||
if !jitterPlaybackReady(false, 40*time.Millisecond, 40*time.Millisecond) {
|
||||
t.Fatal("jitter playback did not start after initial buffer filled")
|
||||
}
|
||||
if !jitterPlaybackReady(true, 1) {
|
||||
if !jitterPlaybackReady(true, 0, 40*time.Millisecond) {
|
||||
t.Fatal("jitter playback paused while refilling after startup")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAudioPacketDurationUsesStereoFrameCount(t *testing.T) {
|
||||
packet := &gumble.AudioPacket{AudioBuffer: make(gumble.AudioBuffer, 2*gumble.AudioDefaultFrameSize)}
|
||||
if got := audioPacketDuration(packet); got != 10*time.Millisecond {
|
||||
t.Fatalf("audioPacketDuration = %v, want 10ms", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderRejectsWorkAfterShutdown(t *testing.T) {
|
||||
s := &Stream{renderClosed: true}
|
||||
called := false
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
barnlog "git.stormux.org/storm/barnard/log"
|
||||
|
||||
@@ -84,6 +85,8 @@ func main() {
|
||||
configSet := false
|
||||
certificateSet := false
|
||||
buffers := flag.Int("buffers", 16, "number of audio buffers to use")
|
||||
audioInterval := flag.Int("audio-interval", 10, "outgoing audio packet duration in ms (10, 20, 40, or 60)")
|
||||
jitterBuffer := flag.Int("jitter-buffer", 40, "incoming per-user audio buffer in ms (0, 20, 40, or 60)")
|
||||
profile := flag.Bool("profile", false, "add http server to serve profiles")
|
||||
noiseSuppressionEnabled := flag.Bool("noise-suppression", false, "enable noise suppression for microphone input")
|
||||
autoTransmit := flag.Bool("auto-transmit", false, "start transmitting immediately on connect")
|
||||
@@ -94,6 +97,14 @@ func main() {
|
||||
logFile := flag.String("logfile", "", "write logs to this file (logging is disabled when omitted)")
|
||||
|
||||
flag.Parse()
|
||||
selectedAudioInterval, err := audioIntervalDuration(*audioInterval)
|
||||
if err != nil {
|
||||
handle_raw_error(err)
|
||||
}
|
||||
selectedJitterBuffer, err := jitterBufferDuration(*jitterBuffer)
|
||||
if err != nil {
|
||||
handle_raw_error(err)
|
||||
}
|
||||
|
||||
// Set up logging
|
||||
var level barnlog.Level
|
||||
@@ -192,6 +203,8 @@ func main() {
|
||||
NoiseSuppressor: noise.NewSuppressor(),
|
||||
}
|
||||
b.Config.Buffers = *buffers
|
||||
b.Config.AudioInterval = selectedAudioInterval
|
||||
b.Config.IncomingAudioBuffer = selectedJitterBuffer
|
||||
b.Config.DisableUDP = *tcpOnly
|
||||
|
||||
b.Hotkeys = b.UserConfig.GetHotkeys()
|
||||
@@ -238,6 +251,30 @@ func main() {
|
||||
handle_error(&b)
|
||||
}
|
||||
|
||||
// audioIntervalDuration converts the packet duration requested at startup to
|
||||
// one of the Opus durations supported by Mumble.
|
||||
func audioIntervalDuration(milliseconds int) (time.Duration, error) {
|
||||
interval := time.Duration(milliseconds) * time.Millisecond
|
||||
switch interval {
|
||||
case 10 * time.Millisecond, 20 * time.Millisecond, 40 * time.Millisecond, 60 * time.Millisecond:
|
||||
return interval, nil
|
||||
default:
|
||||
return 0, fmt.Errorf("audio interval must be 10, 20, 40, or 60 ms, got %d", milliseconds)
|
||||
}
|
||||
}
|
||||
|
||||
// jitterBufferDuration converts the requested incoming playout delay to a
|
||||
// supported duration. Zero starts playback without an initial safety buffer.
|
||||
func jitterBufferDuration(milliseconds int) (time.Duration, error) {
|
||||
interval := time.Duration(milliseconds) * time.Millisecond
|
||||
switch interval {
|
||||
case 0, 20 * time.Millisecond, 40 * time.Millisecond, 60 * time.Millisecond:
|
||||
return interval, nil
|
||||
default:
|
||||
return 0, fmt.Errorf("jitter buffer must be 0, 20, 40, or 60 ms, got %d", milliseconds)
|
||||
}
|
||||
}
|
||||
|
||||
// serverAddress adds Mumble's default port without corrupting an IPv6 literal.
|
||||
func serverAddress(address string) string {
|
||||
if _, port, err := net.SplitHostPort(address); err == nil && port != "" {
|
||||
|
||||
Reference in New Issue
Block a user