Synchronize audio stream teardown
This commit is contained in:
committed by
Brandon McGinty
parent
e8e2fc46a9
commit
2eed3c809f
+22
-11
@@ -256,6 +256,7 @@ func (c *Client) dispatchAudio(user *User, packet *AudioPacket) {
|
|||||||
listeners := &c.Config.AudioListeners
|
listeners := &c.Config.AudioListeners
|
||||||
listeners.mu.Lock()
|
listeners.mu.Lock()
|
||||||
type delivery struct {
|
type delivery struct {
|
||||||
|
item *audioEventItem
|
||||||
listener AudioListener
|
listener AudioListener
|
||||||
ch chan *AudioPacket
|
ch chan *AudioPacket
|
||||||
new bool
|
new bool
|
||||||
@@ -272,7 +273,7 @@ func (c *Client) dispatchAudio(user *User, packet *AudioPacket) {
|
|||||||
ch = make(chan *AudioPacket, bufferSize)
|
ch = make(chan *AudioPacket, bufferSize)
|
||||||
item.streams[user] = ch
|
item.streams[user] = ch
|
||||||
}
|
}
|
||||||
deliveries = append(deliveries, delivery{item.listener, ch, newStream})
|
deliveries = append(deliveries, delivery{item, item.listener, ch, newStream})
|
||||||
}
|
}
|
||||||
listeners.mu.Unlock()
|
listeners.mu.Unlock()
|
||||||
|
|
||||||
@@ -281,12 +282,20 @@ func (c *Client) dispatchAudio(user *User, packet *AudioPacket) {
|
|||||||
log.Debug("new audio stream from %s (session=%d)", user.Name, user.Session)
|
log.Debug("new audio stream from %s (session=%d)", user.Name, user.Session)
|
||||||
delivery.listener.OnAudioStream(&AudioStreamEvent{Client: c, User: user, C: delivery.ch})
|
delivery.listener.OnAudioStream(&AudioStreamEvent{Client: c, User: user, C: delivery.ch})
|
||||||
}
|
}
|
||||||
select {
|
// User removal can run on a different protocol goroutine. Keep the
|
||||||
case delivery.ch <- packet:
|
// listener lock while sending so it cannot close this stream between
|
||||||
default:
|
// the active-stream check and the channel send.
|
||||||
// Never allow a slow listener to block protocol processing.
|
listeners.mu.Lock()
|
||||||
log.Debug("dropping buffered audio for slow listener (session=%d)", user.Session)
|
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.
|
||||||
|
log.Debug("dropping buffered audio for slow listener (session=%d)", user.Session)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
listeners.mu.Unlock()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -578,16 +587,18 @@ func (c *Client) handleUserRemove(buffer []byte) error {
|
|||||||
}
|
}
|
||||||
delete(c.Users, session)
|
delete(c.Users, session)
|
||||||
|
|
||||||
// Close audio stream channels for the disconnected user.
|
// Close audio stream channels for the disconnected user. UDP audio may
|
||||||
// This is safe because handleUserRemove and handleUDPTunnel
|
// still be dispatched concurrently, so the audio-listener lock also
|
||||||
// both run in the serialized readRoutine; no concurrent send
|
// protects its stream maps and channel sends.
|
||||||
// on these channels is possible once the user is removed.
|
listeners := &c.Config.AudioListeners
|
||||||
for item := c.Config.AudioListeners.head; item != nil; item = item.next {
|
listeners.mu.Lock()
|
||||||
|
for item := listeners.head; item != nil; item = item.next {
|
||||||
if ch, ok := item.streams[event.User]; ok {
|
if ch, ok := item.streams[event.User]; ok {
|
||||||
close(ch)
|
close(ch)
|
||||||
delete(item.streams, event.User)
|
delete(item.streams, event.User)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
listeners.mu.Unlock()
|
||||||
|
|
||||||
if packet.Reason != nil {
|
if packet.Reason != nil {
|
||||||
event.String = *packet.Reason
|
event.String = *packet.Reason
|
||||||
|
|||||||
Reference in New Issue
Block a user