improved delivery tests
This commit is contained in:
4
async.go
4
async.go
@@ -69,7 +69,6 @@ func (t *AsyncTopic[T]) run() {
|
||||
select {
|
||||
case newCallback, more := <-t.subscribeCh:
|
||||
if !more {
|
||||
// Ignore subscribeCh close. The publishCh will dictate when to exit this loop.
|
||||
drainedSubscribe = true
|
||||
break
|
||||
}
|
||||
@@ -79,7 +78,6 @@ func (t *AsyncTopic[T]) run() {
|
||||
|
||||
case msg, more := <-t.publishCh:
|
||||
if !more {
|
||||
// No more published messages, promise was fulfilled and we can return
|
||||
drainedPublish = true
|
||||
break
|
||||
}
|
||||
@@ -123,7 +121,7 @@ func (t *AsyncTopic[T]) Subscribe(fn Subscriber[T]) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Feed allows you to for/range to consume future published messages. The supporting subscriber will eventually be discarded after you exit the for loop.
|
||||
// Feed allows the usage of for/range to consume future published messages. The supporting subscriber will eventually be discarded after you exit the for loop.
|
||||
func (t *AsyncTopic[T]) Feed() iter.Seq[T] {
|
||||
feed := make(chan T, 1)
|
||||
done := make(chan struct{})
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
package gubgub
|
||||
|
||||
// sequentialDelivery delivers a message to each subscriber sequentially. For performance reasons
|
||||
// this might mutate the subscribers slice inplace. Please overwrite it with the result of this
|
||||
// sequentialDelivery effentiently delivers a message to each subscriber sequentially. For
|
||||
// performance reasons this might mutate the subscribers slice inplace. Please overwrite it with
|
||||
// the result of this
|
||||
// call.
|
||||
func sequentialDelivery[T any](msg T, subscribers []Subscriber[T]) []Subscriber[T] {
|
||||
last := len(subscribers) - 1
|
||||
|
||||
@@ -12,18 +12,18 @@ func TestSequentialDelivery(t *testing.T) {
|
||||
feedback := make([]int, 0, 3)
|
||||
|
||||
subscribers := []Subscriber[int]{
|
||||
func(i int) bool {
|
||||
assert.Equalf(t, testMsg, i, "expected %d but got %d", testMsg, i)
|
||||
func(msg int) bool {
|
||||
assert.Equalf(t, testMsg, msg, "expected %d but got %d", testMsg, msg)
|
||||
feedback = append(feedback, 1)
|
||||
return true
|
||||
},
|
||||
func(i int) bool {
|
||||
assert.Equalf(t, testMsg, i, "expected %d but got %d", testMsg, i)
|
||||
feedback = append(feedback, 2)
|
||||
return false
|
||||
},
|
||||
func(i int) bool {
|
||||
assert.Equalf(t, testMsg, i, "expected %d but got %d", testMsg, i)
|
||||
func(msg int) bool {
|
||||
assert.Equalf(t, testMsg, msg, "expected %d but got %d", testMsg, msg)
|
||||
feedback = append(feedback, 2)
|
||||
return true
|
||||
},
|
||||
func(msg int) bool {
|
||||
assert.Equalf(t, testMsg, msg, "expected %d but got %d", testMsg, msg)
|
||||
feedback = append(feedback, 3)
|
||||
return true
|
||||
},
|
||||
@@ -38,4 +38,25 @@ func TestSequentialDelivery(t *testing.T) {
|
||||
|
||||
assert.Len(t, finalSubscribers, len(nextSubscribers), "expected to have the same subscribers")
|
||||
assert.Len(t, feedback, 5, "one or more subscriber was not called")
|
||||
|
||||
assertContainsExactlyN(t, 1, 1, feedback)
|
||||
assertContainsExactlyN(t, 2, 2, feedback)
|
||||
assertContainsExactlyN(t, 3, 2, feedback)
|
||||
}
|
||||
|
||||
func assertContainsExactlyN[T comparable](t testing.TB, exp T, n int, slice []T) {
|
||||
t.Helper()
|
||||
|
||||
var found int
|
||||
for _, v := range slice {
|
||||
if exp == v {
|
||||
found++
|
||||
}
|
||||
}
|
||||
|
||||
if n > found {
|
||||
t.Errorf("contains too few '%v': expected %d but found %d", exp, n, found)
|
||||
} else if n < found {
|
||||
t.Errorf("contains too many '%v': expected %d but found %d", exp, n, found)
|
||||
}
|
||||
}
|
||||
|
||||
BIN
gubgub.test
BIN
gubgub.test
Binary file not shown.
15
sync.go
15
sync.go
@@ -2,8 +2,8 @@ package gubgub
|
||||
|
||||
import "sync"
|
||||
|
||||
// SyncTopic allows any message T to be broadcast to subscribers. Publishing and Subscribing
|
||||
// happens synchronously (block).
|
||||
// SyncTopic is the simplest and most naive topic. It allows any message T to be broadcast to
|
||||
// subscribers. Publishing and Subscribing happens synchronously (block).
|
||||
type SyncTopic[T any] struct {
|
||||
mu sync.Mutex
|
||||
subscribers []Subscriber[T]
|
||||
@@ -19,16 +19,7 @@ func (t *SyncTopic[T]) Publish(msg T) error {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
|
||||
keepers := make([]Subscriber[T], 0, len(t.subscribers))
|
||||
|
||||
for _, callback := range t.subscribers {
|
||||
keep := callback(msg)
|
||||
if keep {
|
||||
keepers = append(keepers, callback)
|
||||
}
|
||||
}
|
||||
|
||||
t.subscribers = keepers
|
||||
t.subscribers = sequentialDelivery(msg, t.subscribers)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user