Synchronize listener detach and delivery
This commit is contained in:
committed by
Brandon McGinty
parent
e309d22137
commit
9f90bd5c66
@@ -129,7 +129,7 @@ Priority 1: transport, lifecycle, and correctness
|
||||
channel. Also synchronize Volume, which is currently read and written
|
||||
without protection.
|
||||
|
||||
16. Audio listener/event listener detach is not concurrency-safe
|
||||
[x] 16. Audio listener/event listener detach is not concurrency-safe
|
||||
Files: gumble/gumble/listeners.go, gumble/gumble/audiolisteners.go
|
||||
Event listener detach has no lock; audio detach removes streams without
|
||||
closing/joining them. Concurrent attach/detach/delivery can corrupt linked
|
||||
|
||||
@@ -7,11 +7,16 @@ type audioEventItem struct {
|
||||
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
|
||||
if e.prev == nil {
|
||||
e.parent.head = e.next
|
||||
} else {
|
||||
@@ -22,6 +27,7 @@ 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
|
||||
|
||||
+39
-97
@@ -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) })
|
||||
}
|
||||
|
||||
@@ -0,0 +1,38 @@
|
||||
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)
|
||||
item.Detach()
|
||||
item.Detach()
|
||||
listeners.mu.Lock()
|
||||
defer listeners.mu.Unlock()
|
||||
if listeners.head != nil || listeners.tail != nil {
|
||||
t.Fatal("detached audio listener remains linked")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user