make the Resub unsub channel buffered

This commit is contained in:
inphi 2023-10-18 15:55:58 -04:00
parent 4964f1953e
commit ad2fa66e91
No known key found for this signature in database
GPG key ID: B61066A1A33F5D24

View file

@ -120,7 +120,7 @@ func ResubscribeErr(backoffMax time.Duration, fn ResubscribeErrFunc) Subscriptio
backoffMax: backoffMax, backoffMax: backoffMax,
fn: fn, fn: fn,
err: make(chan error), err: make(chan error),
unsub: make(chan struct{}), unsub: make(chan struct{}, 1),
} }
go s.loop() go s.loop()
return s return s
@ -154,20 +154,13 @@ func (s *resubscribeSub) Err() <-chan error {
func (s *resubscribeSub) loop() { func (s *resubscribeSub) loop() {
var done bool var done bool
var unsubbed bool defer close(s.err)
defer func() {
close(s.err)
// Read unsub chan to avoid blocking Unsubscribe
if !unsubbed {
<-s.unsub
}
}()
for !done { for !done {
sub := s.subscribe() sub := s.subscribe()
if sub == nil { if sub == nil {
break break
} }
done, unsubbed = s.waitForError(sub) done = s.waitForError(sub)
sub.Unsubscribe() sub.Unsubscribe()
} }
} }
@ -204,14 +197,14 @@ func (s *resubscribeSub) subscribe() Subscription {
} }
} }
func (s *resubscribeSub) waitForError(sub Subscription) (bool, bool) { func (s *resubscribeSub) waitForError(sub Subscription) bool {
defer sub.Unsubscribe() defer sub.Unsubscribe()
select { select {
case err := <-sub.Err(): case err := <-sub.Err():
s.lastSubErr = err s.lastSubErr = err
return err == nil, false return err == nil
case <-s.unsub: case <-s.unsub:
return true, true return true
} }
} }