From 28b026c90fe5ab23e0addb7e51d2c4c445214d72 Mon Sep 17 00:00:00 2001 From: Brandon McGinty Date: Thu, 20 Aug 2026 14:20:11 -0400 Subject: [PATCH] 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 --- gumble/gumble/audiolisteners.go | 17 +++ gumble/gumble/client.go | 9 +- .../gumble/contextaction_regression_test.go | 52 +++++++ gumble/gumble/contextactions.go | 5 +- gumble/gumble/handlers.go | 84 ++++++++--- gumble/gumble/handlers_regression_test.go | 28 ++++ gumble/gumble/listeners.go | 136 +++++------------- gumble/gumble/listeners_regression_test.go | 43 ++++++ 8 files changed, 250 insertions(+), 124 deletions(-) create mode 100644 gumble/gumble/contextaction_regression_test.go create mode 100644 gumble/gumble/handlers_regression_test.go create mode 100644 gumble/gumble/listeners_regression_test.go diff --git a/gumble/gumble/audiolisteners.go b/gumble/gumble/audiolisteners.go index 7bf0e80..e56a4ea 100644 --- a/gumble/gumble/audiolisteners.go +++ b/gumble/gumble/audiolisteners.go @@ -1,13 +1,26 @@ package gumble +import "sync" + type audioEventItem struct { parent *AudioListeners prev, next *audioEventItem listener AudioListener streams map[*User]chan *AudioPacket + detached bool } 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 { e.parent.head = e.next } else { @@ -18,16 +31,20 @@ func (e *audioEventItem) Detach() { } else { e.next.prev = e.prev } + e.prev, e.next = nil, nil } // AudioListeners is a list of audio listeners. Each attached listener is // called in sequence when a new user audio stream begins. type AudioListeners struct { + mu sync.Mutex head, tail *audioEventItem } // Attach adds a new audio listener to the end of the current list of listeners. func (e *AudioListeners) Attach(listener AudioListener) Detacher { + e.mu.Lock() + defer e.mu.Unlock() item := &audioEventItem{ parent: e, prev: e.tail, diff --git a/gumble/gumble/client.go b/gumble/gumble/client.go index e77395b..3ad6fcc 100644 --- a/gumble/gumble/client.go +++ b/gumble/gumble/client.go @@ -102,10 +102,11 @@ func DialWithDialer(dialer *net.Dialer, config *Config, tlsConfig *tls.Config) ( } client := &Client{ - Conn: NewConn(conn), - Config: config, - Users: make(Users), - Channels: make(Channels), + Conn: NewConn(conn), + Config: config, + Users: make(Users), + Channels: make(Channels), + ContextActions: make(ContextActions), permissions: make(map[uint32]*Permission), diff --git a/gumble/gumble/contextaction_regression_test.go b/gumble/gumble/contextaction_regression_test.go new file mode 100644 index 0000000..021a39c --- /dev/null +++ b/gumble/gumble/contextaction_regression_test.go @@ -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) + } +} diff --git a/gumble/gumble/contextactions.go b/gumble/gumble/contextactions.go index 6dd0c16..ee58bfd 100644 --- a/gumble/gumble/contextactions.go +++ b/gumble/gumble/contextactions.go @@ -3,9 +3,10 @@ package gumble // ContextActions is a map of ContextActions. type ContextActions map[string]*ContextAction -func (c ContextActions) create(action string) *ContextAction { +func (c ContextActions) create(client *Client, action string) *ContextAction { contextAction := &ContextAction{ - Name: action, + Name: action, + client: client, } c[action] = contextAction return contextAction diff --git a/gumble/gumble/handlers.go b/gumble/gumble/handlers.go index a0278da..b9775f6 100644 --- a/gumble/gumble/handlers.go +++ b/gumble/gumble/handlers.go @@ -156,28 +156,56 @@ func (c *Client) handleUDPTunnel(buffer []byte) error { event.HasPosition = true } - c.volatile.Lock() - 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() - + c.dispatchAudio(user, &event) 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 { return errUnimplementedHandler } @@ -460,6 +488,20 @@ func (c *Client) handleUserRemove(buffer []byte) error { delete(event.User.Channel.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 { event.String = *packet.Reason } @@ -551,7 +593,7 @@ func (c *Client) handleUserState(buffer []byte) error { } newChannel := c.Channels[*packet.ChannelId] if newChannel == nil { - c.volatile.Lock() + c.volatile.Unlock() return errInvalidProtobuf } if newChannel != user.Channel { @@ -924,7 +966,7 @@ func (c *Client) handleContextActionModify(buffer []byte) error { return nil } event.Type = ContextActionAdd - contextAction := c.ContextActions.create(*packet.Action) + contextAction := c.ContextActions.create(c, *packet.Action) if packet.Text != nil { contextAction.Label = *packet.Text } diff --git a/gumble/gumble/handlers_regression_test.go b/gumble/gumble/handlers_regression_test.go new file mode 100644 index 0000000..f1a86d4 --- /dev/null +++ b/gumble/gumble/handlers_regression_test.go @@ -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") + } +} diff --git a/gumble/gumble/listeners.go b/gumble/gumble/listeners.go index 3d100c0..c9267e6 100644 --- a/gumble/gumble/listeners.go +++ b/gumble/gumble/listeners.go @@ -1,12 +1,21 @@ package gumble +import "sync" + type eventItem struct { parent *Listeners prev, next *eventItem listener EventListener + detached bool } func (e *eventItem) Detach() { + e.parent.mu.Lock() + defer e.parent.mu.Unlock() + if e.detached { + return + } + e.detached = true if e.prev == nil { e.parent.head = e.next } else { @@ -17,21 +26,20 @@ func (e *eventItem) Detach() { } else { e.next.prev = e.prev } + e.prev, e.next = nil, nil } -// Listeners is a list of event listeners. Each attached listener is called in -// sequence when a Client event is triggered. +// Listeners is a list of event listeners. Delivery uses a snapshot, so an +// attach or detach from another goroutine cannot corrupt iteration. type Listeners struct { + mu sync.Mutex 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 { - item := &eventItem{ - parent: e, - prev: e.tail, - listener: listener, - } + e.mu.Lock() + defer e.mu.Unlock() + item := &eventItem{parent: e, prev: e.tail, listener: listener} if e.head == nil { e.head = item } @@ -42,112 +50,46 @@ func (e *Listeners) Attach(listener EventListener) Detacher { 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) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnConnect(event) }) } - func (e *Listeners) onDisconnect(event *DisconnectEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnDisconnect(event) }) } - func (e *Listeners) onTextMessage(event *TextMessageEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnTextMessage(event) }) } - func (e *Listeners) onUserChange(event *UserChangeEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnUserChange(event) }) } - func (e *Listeners) onChannelChange(event *ChannelChangeEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnChannelChange(event) }) } - func (e *Listeners) onPermissionDenied(event *PermissionDeniedEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnPermissionDenied(event) }) } - func (e *Listeners) onUserList(event *UserListEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnUserList(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) onACL(event *ACLEvent) { e.dispatch(func(l EventListener) { l.OnACL(event) }) } func (e *Listeners) onBanList(event *BanListEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnBanList(event) }) } - func (e *Listeners) onContextActionChange(event *ContextActionChangeEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnContextActionChange(event) }) } - func (e *Listeners) onServerConfig(event *ServerConfigEvent) { - event.Client.volatile.Lock() - 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() + e.dispatch(func(l EventListener) { l.OnServerConfig(event) }) } diff --git a/gumble/gumble/listeners_regression_test.go b/gumble/gumble/listeners_regression_test.go new file mode 100644 index 0000000..c26a024 --- /dev/null +++ b/gumble/gumble/listeners_regression_test.go @@ -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") + } +}