Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions ringbuf_overflow_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,11 @@ func (d *Data) String() string {
return fmt.Sprintf("Data{ID: %v}", d.ID)
}

// Test extremely rare case of buffer write position overflow (after 2^64 writes).
//
// - At 100M ops/sec: overflow occurs after ~5.8 years
// - At 1B ops/sec: overflow occurs after ~584 days
// - At 10B ops/sec: overflow occurs after ~58 days
func TestWritePosOverflow(t *testing.T) {
stream := New[*Data](100)

Expand Down
84 changes: 25 additions & 59 deletions ringbuf_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,49 +21,7 @@ func (d *Data) String() string {
return fmt.Sprintf("Data{ID: %v, Name: %v}", d.ID, d.Name)
}

func TestBasic(t *testing.T) {
stream := ringbuf.New[*Data](100)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

sub1 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub1"})
sub2 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub2"})
sub3 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub3"})

wg := sync.WaitGroup{}
wg.Add(3)
for _, sub := range []*ringbuf.Subscriber[*Data]{sub1, sub2, sub3} {
go func() {
sub := sub
defer wg.Done()

for val := range sub.Iter() {
t.Logf("%v: Reading %+v", sub.Name, val)
}
if err := sub.Err(); !errors.Is(err, context.Canceled) {
t.Errorf("%v: %v", sub.Name, err)
}
}()
}

for i := range 1000 {
v := &Data{ID: i, Name: fmt.Sprintf("%v", i)}
t.Logf("writer: Writing %+v", v)
stream.Write(v)
time.Sleep(100 * time.Microsecond)
}

cancel() // Terminate the readers.

last := &Data{ID: 1001, Name: "last"}
t.Logf("writer: Writing %+v", last)
stream.Write(last)

wg.Wait()
}

func TestRingBuf(t *testing.T) {
func TestRingbuf(t *testing.T) {
bufferSize := uint64(2_000)
numItems := 10_000
numReaders := 2_000
Expand All @@ -86,30 +44,38 @@ func TestRingBuf(t *testing.T) {
defer wg.Done()
sub := sub

items := make([]*Data, 64)
var count int
for val := range sub.Iter() {
if val.ID != count {
t.Errorf("unexpected data: expected %v, got %v", count, val)
cancel()
return
for {
n, err := sub.Read(items)
if err != nil {
if !errors.Is(err, io.EOF) {
t.Errorf("unexpected error: %v", err)
cancel()
return
}
break
}
if val.Name != fmt.Sprintf("%v", count) {
t.Errorf("unexpected data: expected %v, got %v", count, val)
cancel()
return

for i := range n {
val := items[i]
if val.ID != count {
t.Errorf("unexpected data: expected %v, got %v", count, val)
cancel()
return
}
if val.Name != fmt.Sprintf("%v", count) {
t.Errorf("unexpected data: expected %v, got %v", count, val)
cancel()
return
}
count++
}
count++
}

if count != numItems {
t.Errorf("expected %v items, got %v", numItems, count)
}

if err := sub.Err(); !errors.Is(err, io.EOF) {
t.Errorf("unexpected error: %v", err)
cancel()
return
}
}()
}

Expand Down
6 changes: 2 additions & 4 deletions subscriber_drain_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,8 @@ import (
"time"
)

// Regression test for drain-after-close semantics:
// Read must not return EOF while there are still buffered items available to drain,
// even if Close races with the reader's observation of writePos.
func TestReadDrainsAfterClose(t *testing.T) {
// Read() must not return io.EOF until the data is drained, even if Close() was called.
func TestSubscriberDrainsClosedStream(t *testing.T) {
const iters = 2000

for i := 0; i < iters; i++ {
Expand Down
55 changes: 55 additions & 0 deletions subscriber_iter_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
package ringbuf_test

import (
"context"
"errors"
"fmt"
"sync"
"testing"
"time"

"github.com/golang-cz/ringbuf"
)

func TestSubscriberIter(t *testing.T) {
stream := ringbuf.New[*Data](100)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

sub1 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub1"})
sub2 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub2"})
sub3 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub3"})

wg := sync.WaitGroup{}
wg.Add(3)
for _, sub := range []*ringbuf.Subscriber[*Data]{sub1, sub2, sub3} {
go func() {
sub := sub
defer wg.Done()

for val := range sub.Iter() {
t.Logf("%v: Reading %+v", sub.Name, val)
}
if err := sub.Err(); !errors.Is(err, context.Canceled) {
t.Errorf("%v: %v", sub.Name, err)
}
}()
}

for i := range 1000 {
v := &Data{ID: i, Name: fmt.Sprintf("%v", i)}
t.Logf("writer: Writing %+v", v)
stream.Write(v)
time.Sleep(100 * time.Microsecond)
}

cancel() // Terminate the readers.

// Wake readers up so they can observe ctx cancellation.
last := &Data{ID: 1001, Name: "last"}
t.Logf("writer: Writing %+v", last)
stream.Write(last)

wg.Wait()
}
63 changes: 63 additions & 0 deletions subscriber_read_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
package ringbuf_test

import (
"context"
"errors"
"fmt"
"sync"
"testing"
"time"

"github.com/golang-cz/ringbuf"
)

func TestSubscriberRead(t *testing.T) {
stream := ringbuf.New[*Data](100)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

sub1 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub1"})
sub2 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub2"})
sub3 := stream.Subscribe(ctx, &ringbuf.SubscribeOpts{Name: "sub3"})

wg := sync.WaitGroup{}
wg.Add(3)
for _, sub := range []*ringbuf.Subscriber[*Data]{sub1, sub2, sub3} {
go func() {
sub := sub
defer wg.Done()

items := make([]*Data, 16)
for {
n, err := sub.Read(items)
if err != nil {
if !errors.Is(err, context.Canceled) {
t.Errorf("%v: %v", sub.Name, err)
}
return
}

for i := range n {
t.Logf("%v: Reading %+v", sub.Name, items[i])
}
}
}()
}

for i := range 1000 {
v := &Data{ID: i, Name: fmt.Sprintf("%v", i)}
t.Logf("writer: Writing %+v", v)
stream.Write(v)
time.Sleep(100 * time.Microsecond)
}

cancel() // Terminate the readers.

// Wake readers up so they can observe ctx cancellation.
last := &Data{ID: 1001, Name: "last"}
t.Logf("writer: Writing %+v", last)
stream.Write(last)

wg.Wait()
}
2 changes: 1 addition & 1 deletion seek_example_test.go → subscriber_reconnect_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ type Event struct {
Time time.Time
}

func TestSeek_Reconnect(t *testing.T) {
func TestSubscriberReconnect(t *testing.T) {
t.Parallel()

ctx, cancel := context.WithCancel(context.Background())
Expand Down
2 changes: 1 addition & 1 deletion ringbuf_seek_test.go → subscriber_seek_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import (
"github.com/golang-cz/ringbuf"
)

func TestSeek_Input(t *testing.T) {
func TestSubscriberSeek(t *testing.T) {
cases := []struct {
name string
ids []int
Expand Down
Loading