mirror of
https://github.com/libp2p/go-eventbus.git
synced 2025-04-02 13:40:12 +08:00
Merge pull request #11 from libp2p/fix/things
Fix close deadlock and Sub type error
This commit is contained in:
commit
df5be7d7dd
28
basic.go
28
basic.go
@ -104,9 +104,21 @@ func (s *sub) Out() <-chan interface{} {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (s *sub) Close() error {
|
func (s *sub) Close() error {
|
||||||
close(s.ch)
|
stop := make(chan struct{})
|
||||||
|
go func() {
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-s.ch:
|
||||||
|
case <-stop:
|
||||||
|
close(s.ch)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
for _, n := range s.nodes {
|
for _, n := range s.nodes {
|
||||||
n.lk.Lock()
|
n.lk.Lock()
|
||||||
|
|
||||||
for i := 0; i < len(n.sinks); i++ {
|
for i := 0; i < len(n.sinks); i++ {
|
||||||
if n.sinks[i] == s.ch {
|
if n.sinks[i] == s.ch {
|
||||||
n.sinks[i], n.sinks[len(n.sinks)-1] = n.sinks[len(n.sinks)-1], nil
|
n.sinks[i], n.sinks[len(n.sinks)-1] = n.sinks[len(n.sinks)-1], nil
|
||||||
@ -114,12 +126,16 @@ func (s *sub) Close() error {
|
|||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
tryDrop := len(n.sinks) == 0 && atomic.LoadInt32(&n.nEmitters) == 0
|
tryDrop := len(n.sinks) == 0 && atomic.LoadInt32(&n.nEmitters) == 0
|
||||||
|
|
||||||
n.lk.Unlock()
|
n.lk.Unlock()
|
||||||
|
|
||||||
if tryDrop {
|
if tryDrop {
|
||||||
s.dropper(n.typ)
|
s.dropper(n.typ)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
close(stop)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -148,12 +164,14 @@ func (b *basicBus) Subscribe(evtTypes interface{}, opts ...event.SubscriptionOpt
|
|||||||
dropper: b.tryDropNode,
|
dropper: b.tryDropNode,
|
||||||
}
|
}
|
||||||
|
|
||||||
for i, etyp := range types {
|
for _, etyp := range types {
|
||||||
typ := reflect.TypeOf(etyp)
|
if reflect.TypeOf(etyp).Kind() != reflect.Ptr {
|
||||||
|
|
||||||
if typ.Kind() != reflect.Ptr {
|
|
||||||
return nil, errors.New("subscribe called with non-pointer type")
|
return nil, errors.New("subscribe called with non-pointer type")
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for i, etyp := range types {
|
||||||
|
typ := reflect.TypeOf(etyp)
|
||||||
|
|
||||||
err = b.withNode(typ.Elem(), func(n *node) {
|
err = b.withNode(typ.Elem(), func(n *node) {
|
||||||
n.sinks = append(n.sinks, out.ch)
|
n.sinks = append(n.sinks, out.ch)
|
||||||
|
@ -308,6 +308,49 @@ func TestStateful(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestCloseBlocking(t *testing.T) {
|
||||||
|
bus := NewBus()
|
||||||
|
em, err := bus.Emitter(new(EventB))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
sub, err := bus.Subscribe(new(EventB))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
em.Emit(EventB(159))
|
||||||
|
}()
|
||||||
|
|
||||||
|
time.Sleep(10 * time.Millisecond) // make sure that emit is blocked
|
||||||
|
|
||||||
|
sub.Close() // cancel sub
|
||||||
|
}
|
||||||
|
|
||||||
|
func panicOnTimeout(d time.Duration) {
|
||||||
|
<-time.After(d)
|
||||||
|
panic("timeout reached")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSubFailFully(t *testing.T) {
|
||||||
|
bus := NewBus()
|
||||||
|
em, err := bus.Emitter(new(EventB))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = bus.Subscribe([]interface{}{new(EventB), 5})
|
||||||
|
if err == nil || err.Error() != "subscribe called with non-pointer type" {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
go panicOnTimeout(5 * time.Second)
|
||||||
|
|
||||||
|
em.Emit(EventB(159)) // will hang if sub doesn't fail properly
|
||||||
|
}
|
||||||
|
|
||||||
func testMany(t testing.TB, subs, emits, msgs int, stateful bool) {
|
func testMany(t testing.TB, subs, emits, msgs int, stateful bool) {
|
||||||
if race.WithRace() && subs+emits > 5000 {
|
if race.WithRace() && subs+emits > 5000 {
|
||||||
t.SkipNow()
|
t.SkipNow()
|
||||||
|
Loading…
Reference in New Issue
Block a user