fix: serialize publishing

Ensure that all subscribers see events in the same order. This also ensures that
the subscribers never see the initial "latest" event after some other event.

fixes #16
This commit is contained in:
Steven Allen
2019-06-27 22:33:53 +02:00
parent 04058af20a
commit 25d54bbbec

View File

@@ -173,7 +173,7 @@ func (b *basicBus) Subscribe(evtTypes interface{}, opts ...event.SubscriptionOpt
out.nodes[i] = n out.nodes[i] = n
}, func(n *node) { }, func(n *node) {
if n.keepLast { if n.keepLast {
l := n.last.Load() l := n.last
if l == nil { if l == nil {
return return
} }
@@ -223,7 +223,7 @@ func (b *basicBus) Emitter(evtType interface{}, opts ...event.EmitterOpt) (e eve
type node struct { type node struct {
// Note: make sure to NEVER lock basicBus.lk when this lock is held // Note: make sure to NEVER lock basicBus.lk when this lock is held
lk sync.RWMutex lk sync.Mutex
typ reflect.Type typ reflect.Type
@@ -231,7 +231,7 @@ type node struct {
nEmitters int32 nEmitters int32
keepLast bool keepLast bool
last atomic.Value last interface{}
sinks []chan interface{} sinks []chan interface{}
} }
@@ -248,13 +248,13 @@ func (n *node) emit(event interface{}) {
panic(fmt.Sprintf("Emit called with wrong type. expected: %s, got: %s", n.typ, typ)) panic(fmt.Sprintf("Emit called with wrong type. expected: %s, got: %s", n.typ, typ))
} }
n.lk.RLock() n.lk.Lock()
if n.keepLast { if n.keepLast {
n.last.Store(event) n.last = event
} }
for _, ch := range n.sinks { for _, ch := range n.sinks {
ch <- event ch <- event
} }
n.lk.RUnlock() n.lk.Unlock()
} }