Stop ring consumers when deleting network

This commit is contained in:
Simon Ser 2020-04-01 15:48:56 +02:00
parent 96039320b6
commit 10da094259
No known key found for this signature in database
GPG Key ID: 0FDE7BE0E88F5E48
3 changed files with 27 additions and 1 deletions

View File

@ -739,7 +739,12 @@ func (dc *downstreamConn) runNetwork(net *network, loadHistory bool) {
for {
var closed bool
select {
case <-ch:
case _, ok := <-ch:
if !ok {
closed = true
break
}
uc := net.upstream()
if uc == nil {
dc.logger.Printf("ignoring messages for upstream %q: upstream is disconnected", net.Addr)

20
ring.go
View File

@ -16,6 +16,7 @@ type Ring struct {
lock sync.Mutex
cur uint64
consumers []*RingConsumer
closed bool
}
// NewRing creates a new ring buffer.
@ -31,6 +32,10 @@ func (r *Ring) Produce(msg *irc.Message) {
r.lock.Lock()
defer r.lock.Unlock()
if r.closed {
panic("soju: Ring.Produce called after Close")
}
i := int(r.cur % r.cap)
r.buffer[i] = msg
r.cur++
@ -45,6 +50,21 @@ func (r *Ring) Produce(msg *irc.Message) {
}
}
func (r *Ring) Close() {
r.lock.Lock()
defer r.lock.Unlock()
if r.closed {
panic("soju: Ring.Close called twice")
}
for _, rc := range r.consumers {
close(rc.ch)
}
r.closed = true
}
// NewConsumer creates a new ring buffer consumer.
//
// If seq is nil, the consumer will get messages starting from the last

View File

@ -305,6 +305,7 @@ func (u *user) deleteNetwork(id int64) error {
})
net.Stop()
net.ring.Close()
u.networks = append(u.networks[:i], u.networks[i+1:]...)
return nil
}