diff --git a/core/collection/fifo.go b/core/collection/fifo.go index c310af102..c1e9bc8da 100644 --- a/core/collection/fifo.go +++ b/core/collection/fifo.go @@ -2,11 +2,12 @@ package collection import "sync" +const queueGrowThreshold = 256 + // A Queue is a FIFO queue. type Queue struct { lock sync.Mutex elements []any - size int head int tail int count int @@ -14,9 +15,12 @@ type Queue struct { // NewQueue returns a Queue object. func NewQueue(size int) *Queue { + if size < 1 { + panic("size must be greater than 0") + } + return &Queue{ elements: make([]any, size), - size: size, } } @@ -34,12 +38,12 @@ func (q *Queue) Put(element any) { q.lock.Lock() defer q.lock.Unlock() - if q.head == q.tail && q.count > 0 { - nodes := make([]any, len(q.elements)+q.size) - copy(nodes, q.elements[q.head:]) - copy(nodes[len(q.elements)-q.head:], q.elements[:q.head]) + if q.count == len(q.elements) { + nodes := make([]any, nextQueueCapacity(len(q.elements))) + n := copy(nodes, q.elements[q.head:]) + copy(nodes[n:], q.elements[:q.head]) q.head = 0 - q.tail = len(q.elements) + q.tail = q.count q.elements = nodes } @@ -58,8 +62,19 @@ func (q *Queue) Take() (any, bool) { } element := q.elements[q.head] + q.elements[q.head] = nil q.head = (q.head + 1) % len(q.elements) q.count-- return element, true } + +func nextQueueCapacity(capacity int) int { + if capacity < queueGrowThreshold { + return capacity << 1 + } + + // Use a growth curve similar to Go slices: double small queues, then + // transition smoothly toward 1.25x growth for larger queues. + return capacity + ((capacity + 3*queueGrowThreshold) >> 2) +} diff --git a/core/collection/fifo_test.go b/core/collection/fifo_test.go index 6f1acfdf0..41f19e1ef 100644 --- a/core/collection/fifo_test.go +++ b/core/collection/fifo_test.go @@ -6,96 +6,263 @@ import ( "github.com/stretchr/testify/assert" ) -func TestFifo(t *testing.T) { - elements := [][]byte{ - []byte("hello"), - []byte("world"), - []byte("again"), - } - queue := NewQueue(8) - for i := range elements { - queue.Put(elements[i]) +func TestQueueOrder(t *testing.T) { + tests := []struct { + name string + size int + initial []int + takeBefore int + additional []int + wantCapacity int + }{ + { + name: "within initial capacity", + size: 8, + initial: []int{1, 2, 3}, + wantCapacity: 8, + }, + { + name: "grow from beginning", + size: 2, + initial: []int{1, 2, 3}, + wantCapacity: 4, + }, + { + name: "grow after wrapping", + size: 4, + initial: []int{1, 2, 3, 4}, + takeBefore: 1, + additional: []int{5, 6}, + wantCapacity: 8, + }, + { + name: "grow repeatedly", + size: 1, + initial: sequence(20), + wantCapacity: 32, + }, + { + name: "grow above threshold", + size: queueGrowThreshold, + initial: sequence(queueGrowThreshold + 1), + wantCapacity: queueGrowThreshold * 2, + }, + { + name: "grow well above threshold", + size: 1024, + initial: sequence(1025), + wantCapacity: 1472, + }, } - for _, element := range elements { - body, ok := queue.Take() - assert.True(t, ok) - assert.Equal(t, string(element), string(body.([]byte))) + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + queue := NewQueue(test.size) + assert.True(t, queue.Empty()) + + for _, value := range test.initial { + queue.Put(value) + } + for i := 0; i < test.takeBefore; i++ { + value, ok := queue.Take() + assert.True(t, ok) + assert.Equal(t, test.initial[i], value) + } + for _, value := range test.additional { + queue.Put(value) + } + + want := append([]int(nil), test.initial[test.takeBefore:]...) + want = append(want, test.additional...) + for _, expected := range want { + actual, ok := queue.Take() + assert.True(t, ok) + assert.Equal(t, expected, actual) + } + + assert.Equal(t, test.wantCapacity, len(queue.elements)) + assert.True(t, queue.Empty()) + _, ok := queue.Take() + assert.False(t, ok) + }) } } -func TestTakeTooMany(t *testing.T) { - elements := [][]byte{ - []byte("hello"), - []byte("world"), - []byte("again"), - } - queue := NewQueue(8) - for i := range elements { - queue.Put(elements[i]) +func TestQueueTakeClearsElement(t *testing.T) { + tests := []struct { + name string + size int + operations string + }{ + { + name: "take from beginning", + size: 2, + operations: "ppt", + }, + { + name: "take after wrapping", + size: 2, + operations: "pptptt", + }, + { + name: "take after growing", + size: 2, + operations: "pppt", + }, } - for range elements { - queue.Take() - } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + queue := NewQueue(test.size) + value := 0 - assert.True(t, queue.Empty()) - _, ok := queue.Take() - assert.False(t, ok) -} - -func TestPutMore(t *testing.T) { - elements := [][]byte{ - []byte("hello"), - []byte("world"), - []byte("again"), - } - queue := NewQueue(2) - for i := range elements { - queue.Put(elements[i]) - } - - for _, element := range elements { - body, ok := queue.Take() - assert.True(t, ok) - assert.Equal(t, string(element), string(body.([]byte))) + for _, operation := range test.operations { + switch operation { + case 'p': + value++ + element := new(int) + *element = value + queue.Put(element) + case 't': + index := queue.head + _, ok := queue.Take() + assert.True(t, ok) + assert.Nil(t, queue.elements[index]) + default: + t.Fatalf("unknown operation: %q", operation) + } + } + }) } } -func TestPutMoreWithHeaderNotZero(t *testing.T) { - elements := [][]byte{ - []byte("hello"), - []byte("world"), - []byte("again"), - } - queue := NewQueue(4) - for i := range elements { - queue.Put(elements[i]) +func TestNewQueueWithInvalidSize(t *testing.T) { + tests := []struct { + name string + size int + }{ + { + name: "zero", + }, + { + name: "negative", + size: -1, + }, } - // take 1 - body, ok := queue.Take() - assert.True(t, ok) - element, ok := body.([]byte) - assert.True(t, ok) - assert.Equal(t, element, []byte("hello")) - - // put more - queue.Put([]byte("b4")) - queue.Put([]byte("b5")) // will store in elements[0] - queue.Put([]byte("b6")) // cause expansion - - results := [][]byte{ - []byte("world"), - []byte("again"), - []byte("b4"), - []byte("b5"), - []byte("b6"), - } - - for _, element := range results { - body, ok := queue.Take() - assert.True(t, ok) - assert.Equal(t, string(element), string(body.([]byte))) + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + assert.Panics(t, func() { + NewQueue(test.size) + }) + }) } } + +func BenchmarkQueueGrowth(b *testing.B) { + tests := []struct { + name string + size int + count int + }{ + { + name: "initial_1_count_256", + size: 1, + count: 256, + }, + { + name: "initial_1_count_4096", + size: 1, + count: 4096, + }, + { + name: "initial_1_count_65536", + size: 1, + count: 65536, + }, + { + name: "initial_8_count_4096", + size: 8, + count: 4096, + }, + { + name: "initial_256_count_4096", + size: 256, + count: 4096, + }, + } + + for _, test := range tests { + b.Run(test.name, func(b *testing.B) { + elements := make([]any, test.count) + for i := range elements { + elements[i] = i + } + + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + queue := NewQueue(test.size) + for _, element := range elements { + queue.Put(element) + } + } + }) + } +} + +func BenchmarkQueueWrappedGrowth(b *testing.B) { + tests := []struct { + name string + size int + }{ + { + name: "capacity_8", + size: 8, + }, + { + name: "capacity_256", + size: 256, + }, + { + name: "capacity_4096", + size: 4096, + }, + } + + for _, test := range tests { + b.Run(test.name, func(b *testing.B) { + elements := make([]any, test.size) + for i := range elements { + elements[i] = i + } + + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + queue := NewQueue(test.size) + for _, element := range elements { + queue.Put(element) + } + for j := 0; j < test.size/2; j++ { + queue.Take() + } + for j := 0; j < test.size/2; j++ { + queue.Put(elements[j]) + } + + // The queue is full and wrapped. One more Put triggers growth + // and copies both sides of the ring into FIFO order. + queue.Put(elements[0]) + } + }) + } +} + +func sequence(size int) []int { + values := make([]int, size) + for i := range values { + values[i] = i + } + return values +}