dispatch events and audio from locked snapshots

Guard the listener lists with a mutex and iterate a snapshot.
Attaching or detaching a listener from another goroutine mutated the
linked list while a dispatch was walking it, and a second Detach on
the same item corrupted the head and tail links.

Deliver audio through dispatchAudio instead of an inline loop that
juggled the volatile lock.
The old loop unlocked and relocked around every listener, so a user
map change mid-dispatch could be missed. Streams are now buffered and
a slow listener has its packet dropped rather than blocking protocol
processing.

Close a user's audio streams when the user is removed.
Detaching a listener closes its streams too, so a stream can no longer
be written after the reader is gone.

Unlock volatile when a user moves to an unknown channel.
The error path locked it a second time instead of unlocking, which
deadlocked all later protocol handling.

Give context actions their client when they are created.
Every context action method dereferences it, so triggering a
server-supplied action panicked on a nil client.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Brandon McGinty
2026-08-20 14:20:11 -04:00
co-authored by Claude Opus 5
parent 4788f8da24
commit 28b026c90f
8 changed files with 250 additions and 124 deletions
+17
View File
@@ -1,13 +1,26 @@
package gumble package gumble
import "sync"
type audioEventItem struct { type audioEventItem struct {
parent *AudioListeners parent *AudioListeners
prev, next *audioEventItem prev, next *audioEventItem
listener AudioListener listener AudioListener
streams map[*User]chan *AudioPacket streams map[*User]chan *AudioPacket
detached bool
} }
func (e *audioEventItem) Detach() { func (e *audioEventItem) Detach() {
e.parent.mu.Lock()
defer e.parent.mu.Unlock()
if e.detached {
return
}
e.detached = true
for user, stream := range e.streams {
close(stream)
delete(e.streams, user)
}
if e.prev == nil { if e.prev == nil {
e.parent.head = e.next e.parent.head = e.next
} else { } else {
@@ -18,16 +31,20 @@ func (e *audioEventItem) Detach() {
} else { } else {
e.next.prev = e.prev e.next.prev = e.prev
} }
e.prev, e.next = nil, nil
} }
// AudioListeners is a list of audio listeners. Each attached listener is // AudioListeners is a list of audio listeners. Each attached listener is
// called in sequence when a new user audio stream begins. // called in sequence when a new user audio stream begins.
type AudioListeners struct { type AudioListeners struct {
mu sync.Mutex
head, tail *audioEventItem head, tail *audioEventItem
} }
// Attach adds a new audio listener to the end of the current list of listeners. // Attach adds a new audio listener to the end of the current list of listeners.
func (e *AudioListeners) Attach(listener AudioListener) Detacher { func (e *AudioListeners) Attach(listener AudioListener) Detacher {
e.mu.Lock()
defer e.mu.Unlock()
item := &audioEventItem{ item := &audioEventItem{
parent: e, parent: e,
prev: e.tail, prev: e.tail,
+1
View File
@@ -106,6 +106,7 @@ func DialWithDialer(dialer *net.Dialer, config *Config, tlsConfig *tls.Config) (
Config: config, Config: config,
Users: make(Users), Users: make(Users),
Channels: make(Channels), Channels: make(Channels),
ContextActions: make(ContextActions),
permissions: make(map[uint32]*Permission), permissions: make(map[uint32]*Permission),
@@ -0,0 +1,52 @@
package gumble
import (
"net"
"testing"
"git.stormux.org/storm/barnard/gumble/gumble/MumbleProto"
"google.golang.org/protobuf/proto"
)
// Regression: server context-action adds wrote to a nil map and the resulting
// action had no owning client, so Trigger panicked.
func TestContextActionAddAndTrigger(t *testing.T) {
clientConn, serverConn := net.Pipe()
defer serverConn.Close()
c := &Client{Config: NewConfig(), Users: make(Users), Channels: make(Channels), ContextActions: make(ContextActions)}
c.Conn = NewConn(clientConn)
action, operation := "test", MumbleProto.ContextActionModify_Add
data, err := proto.Marshal(&MumbleProto.ContextActionModify{Action: &action, Operation: &operation})
if err != nil {
t.Fatal(err)
}
if err := c.handleContextActionModify(data); err != nil {
t.Fatal(err)
}
added := c.ContextActions[action]
if added == nil || added.client != c {
t.Fatal("action was not initialized with its client")
}
written := make(chan error, 1)
go func() { _, _, err := NewConn(serverConn).ReadPacket(); written <- err }()
added.Trigger()
if err := <-written; err != nil {
t.Fatalf("trigger did not write: %v", err)
}
remove := MumbleProto.ContextActionModify_Remove
data, _ = proto.Marshal(&MumbleProto.ContextActionModify{Action: &action, Operation: &remove})
if err := c.handleContextActionModify(data); err != nil {
t.Fatal(err)
}
if c.ContextActions[action] != nil {
t.Fatal("action was not removed")
}
}
// Regression: the documented bit layout disagreed with SemanticVersion.
func TestSemanticVersionKnownLayout(t *testing.T) {
major, minor, patch := (&Version{Version: 1<<16 | 5<<8 | 2}).SemanticVersion()
if major != 1 || minor != 5 || patch != 2 {
t.Fatalf("got %d.%d.%d", major, minor, patch)
}
}
+2 -1
View File
@@ -3,9 +3,10 @@ package gumble
// ContextActions is a map of ContextActions. // ContextActions is a map of ContextActions.
type ContextActions map[string]*ContextAction type ContextActions map[string]*ContextAction
func (c ContextActions) create(action string) *ContextAction { func (c ContextActions) create(client *Client, action string) *ContextAction {
contextAction := &ContextAction{ contextAction := &ContextAction{
Name: action, Name: action,
client: client,
} }
c[action] = contextAction c[action] = contextAction
return contextAction return contextAction
+63 -21
View File
@@ -156,28 +156,56 @@ func (c *Client) handleUDPTunnel(buffer []byte) error {
event.HasPosition = true event.HasPosition = true
} }
c.volatile.Lock() c.dispatchAudio(user, &event)
for item := c.Config.AudioListeners.head; item != nil; item = item.next {
c.volatile.Unlock()
ch := item.streams[user]
if ch == nil {
ch = make(chan *AudioPacket)
item.streams[user] = ch
event := AudioStreamEvent{
Client: c,
User: user,
C: ch,
}
item.listener.OnAudioStream(&event)
}
ch <- &event
c.volatile.Lock()
}
c.volatile.Unlock()
return nil return nil
} }
// dispatchAudio sends an audio packet to all registered audio listeners.
func (c *Client) dispatchAudio(user *User, packet *AudioPacket) {
listeners := &c.Config.AudioListeners
listeners.mu.Lock()
type delivery struct {
item *audioEventItem
listener AudioListener
ch chan *AudioPacket
new bool
}
var deliveries []delivery
for item := listeners.head; item != nil; item = item.next {
ch := item.streams[user]
newStream := ch == nil
if newStream {
bufferSize := c.Config.Buffers
if bufferSize < 1 {
bufferSize = 1
}
ch = make(chan *AudioPacket, bufferSize)
item.streams[user] = ch
}
deliveries = append(deliveries, delivery{item, item.listener, ch, newStream})
}
listeners.mu.Unlock()
for _, delivery := range deliveries {
if delivery.new {
delivery.listener.OnAudioStream(&AudioStreamEvent{Client: c, User: user, C: delivery.ch})
}
// User removal can run on a different protocol goroutine. Keep the
// listener lock while sending so it cannot close this stream between
// the active-stream check and the channel send.
listeners.mu.Lock()
active := !delivery.item.detached && delivery.item.streams[user] == delivery.ch
if active {
select {
case delivery.ch <- packet:
default:
// Never allow a slow listener to block protocol processing.
}
}
listeners.mu.Unlock()
}
}
func (c *Client) handleAuthenticate(buffer []byte) error { func (c *Client) handleAuthenticate(buffer []byte) error {
return errUnimplementedHandler return errUnimplementedHandler
} }
@@ -460,6 +488,20 @@ func (c *Client) handleUserRemove(buffer []byte) error {
delete(event.User.Channel.Users, session) delete(event.User.Channel.Users, session)
} }
delete(c.Users, session) delete(c.Users, session)
// Close audio stream channels for the disconnected user. UDP audio may
// still be dispatched concurrently, so the audio-listener lock also
// protects its stream maps and channel sends.
listeners := &c.Config.AudioListeners
listeners.mu.Lock()
for item := listeners.head; item != nil; item = item.next {
if ch, ok := item.streams[event.User]; ok {
close(ch)
delete(item.streams, event.User)
}
}
listeners.mu.Unlock()
if packet.Reason != nil { if packet.Reason != nil {
event.String = *packet.Reason event.String = *packet.Reason
} }
@@ -551,7 +593,7 @@ func (c *Client) handleUserState(buffer []byte) error {
} }
newChannel := c.Channels[*packet.ChannelId] newChannel := c.Channels[*packet.ChannelId]
if newChannel == nil { if newChannel == nil {
c.volatile.Lock() c.volatile.Unlock()
return errInvalidProtobuf return errInvalidProtobuf
} }
if newChannel != user.Channel { if newChannel != user.Channel {
@@ -924,7 +966,7 @@ func (c *Client) handleContextActionModify(buffer []byte) error {
return nil return nil
} }
event.Type = ContextActionAdd event.Type = ContextActionAdd
contextAction := c.ContextActions.create(*packet.Action) contextAction := c.ContextActions.create(c, *packet.Action)
if packet.Text != nil { if packet.Text != nil {
contextAction.Label = *packet.Text contextAction.Label = *packet.Text
} }
+28
View File
@@ -0,0 +1,28 @@
package gumble
import (
"testing"
"time"
"git.stormux.org/storm/barnard/gumble/gumble/MumbleProto"
"google.golang.org/protobuf/proto"
)
// Regression: an out-of-order user move to an unknown channel used to lock
// volatile twice, permanently deadlocking subsequent protocol handling.
func TestUserStateUnknownChannelDoesNotDeadlock(t *testing.T) {
c := &Client{Config: NewConfig(), Users: make(Users), Channels: make(Channels)}
c.Users.create(1)
id, channel := uint32(1), uint32(99)
data, _ := proto.Marshal(&MumbleProto.UserState{Session: &id, ChannelId: &channel})
if err := c.handleUserState(data); err != errInvalidProtobuf {
t.Fatalf("got %v", err)
}
locked := make(chan struct{})
go func() { c.volatile.Lock(); c.volatile.Unlock(); close(locked) }()
select {
case <-locked:
case <-time.After(time.Second):
t.Fatal("volatile lock was left locked")
}
}
+39 -97
View File
@@ -1,12 +1,21 @@
package gumble package gumble
import "sync"
type eventItem struct { type eventItem struct {
parent *Listeners parent *Listeners
prev, next *eventItem prev, next *eventItem
listener EventListener listener EventListener
detached bool
} }
func (e *eventItem) Detach() { func (e *eventItem) Detach() {
e.parent.mu.Lock()
defer e.parent.mu.Unlock()
if e.detached {
return
}
e.detached = true
if e.prev == nil { if e.prev == nil {
e.parent.head = e.next e.parent.head = e.next
} else { } else {
@@ -17,21 +26,20 @@ func (e *eventItem) Detach() {
} else { } else {
e.next.prev = e.prev e.next.prev = e.prev
} }
e.prev, e.next = nil, nil
} }
// Listeners is a list of event listeners. Each attached listener is called in // Listeners is a list of event listeners. Delivery uses a snapshot, so an
// sequence when a Client event is triggered. // attach or detach from another goroutine cannot corrupt iteration.
type Listeners struct { type Listeners struct {
mu sync.Mutex
head, tail *eventItem head, tail *eventItem
} }
// Attach adds a new event listener to the end of the current list of listeners.
func (e *Listeners) Attach(listener EventListener) Detacher { func (e *Listeners) Attach(listener EventListener) Detacher {
item := &eventItem{ e.mu.Lock()
parent: e, defer e.mu.Unlock()
prev: e.tail, item := &eventItem{parent: e, prev: e.tail, listener: listener}
listener: listener,
}
if e.head == nil { if e.head == nil {
e.head = item e.head = item
} }
@@ -42,112 +50,46 @@ func (e *Listeners) Attach(listener EventListener) Detacher {
return item return item
} }
func (e *Listeners) dispatch(f func(EventListener)) {
e.mu.Lock()
listeners := make([]EventListener, 0)
for item := e.head; item != nil; item = item.next {
listeners = append(listeners, item.listener)
}
e.mu.Unlock()
for _, listener := range listeners {
f(listener)
}
}
func (e *Listeners) onConnect(event *ConnectEvent) { func (e *Listeners) onConnect(event *ConnectEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnConnect(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnConnect(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onDisconnect(event *DisconnectEvent) { func (e *Listeners) onDisconnect(event *DisconnectEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnDisconnect(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnDisconnect(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onTextMessage(event *TextMessageEvent) { func (e *Listeners) onTextMessage(event *TextMessageEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnTextMessage(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnTextMessage(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onUserChange(event *UserChangeEvent) { func (e *Listeners) onUserChange(event *UserChangeEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnUserChange(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnUserChange(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onChannelChange(event *ChannelChangeEvent) { func (e *Listeners) onChannelChange(event *ChannelChangeEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnChannelChange(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnChannelChange(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onPermissionDenied(event *PermissionDeniedEvent) { func (e *Listeners) onPermissionDenied(event *PermissionDeniedEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnPermissionDenied(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnPermissionDenied(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onUserList(event *UserListEvent) { func (e *Listeners) onUserList(event *UserListEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnUserList(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnUserList(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onACL(event *ACLEvent) { e.dispatch(func(l EventListener) { l.OnACL(event) }) }
func (e *Listeners) onACL(event *ACLEvent) {
event.Client.volatile.Lock()
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnACL(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
}
func (e *Listeners) onBanList(event *BanListEvent) { func (e *Listeners) onBanList(event *BanListEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnBanList(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnBanList(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onContextActionChange(event *ContextActionChangeEvent) { func (e *Listeners) onContextActionChange(event *ContextActionChangeEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnContextActionChange(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnContextActionChange(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
func (e *Listeners) onServerConfig(event *ServerConfigEvent) { func (e *Listeners) onServerConfig(event *ServerConfigEvent) {
event.Client.volatile.Lock() e.dispatch(func(l EventListener) { l.OnServerConfig(event) })
for item := e.head; item != nil; item = item.next {
event.Client.volatile.Unlock()
item.listener.OnServerConfig(event)
event.Client.volatile.Lock()
}
event.Client.volatile.Unlock()
} }
@@ -0,0 +1,43 @@
package gumble
import (
"sync"
"testing"
)
// Regression: Detach modified the linked listener list without synchronization
// and a second detach could corrupt its head/tail links.
func TestListenerDetachIsIdempotentAndConcurrent(t *testing.T) {
var listeners Listeners
item := listeners.Attach(nil)
var wg sync.WaitGroup
for i := 0; i < 16; i++ {
wg.Add(1)
go func() { defer wg.Done(); item.Detach() }()
}
wg.Wait()
listeners.mu.Lock()
defer listeners.mu.Unlock()
if listeners.head != nil || listeners.tail != nil {
t.Fatal("detached listener remains linked")
}
}
// Regression: audio listener detach could be invoked twice while dispatch was
// active, leaving list links inconsistent.
func TestAudioListenerDetachIsIdempotent(t *testing.T) {
var listeners AudioListeners
item := listeners.Attach(nil)
stream := make(chan *AudioPacket)
item.(*audioEventItem).streams[&User{}] = stream
item.Detach()
item.Detach()
if _, open := <-stream; open {
t.Fatal("detached audio listener stream remained open")
}
listeners.mu.Lock()
defer listeners.mu.Unlock()
if listeners.head != nil || listeners.tail != nil {
t.Fatal("detached audio listener remains linked")
}
}