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") + } +}