Rebuild search indexes asynchronously (#111829)

* Add "debouncer" queue, which can combine incoming elements.

* Rebuild indexes asynchronously.

* Remove duplicate method.

* Fix bleve tests.

* Extracted combineRebuildRequests and added test for it.

* Add TestShouldRebuildIndex

* Added TestFindIndexesForRebuild

* Added TestFindIndexesForRebuild

* Introduce index_rebuild_workers option.

* Add metric for rebuild queue length.

* Add TestRebuildIndexes.

* Fix import.

* Linter, review feedback.
This commit is contained in:
Peter Štibraný
2025-10-01 11:52:09 +02:00
committed by GitHub
parent 91e8eb0e45
commit 707c486a46
11 changed files with 1009 additions and 285 deletions
+119
View File
@@ -0,0 +1,119 @@
package debouncer
import (
"context"
"errors"
"slices"
"sync"
)
type CombineFn[T any] func(a, b T) (c T, ok bool)
// Queue is a queue of elements. Elements added to the queue can be combined together by the provided combiner function.
// Once the queue is closed, no more elements can be added, but Next() will still return remaining elements.
type Queue[T any] struct {
combineFn CombineFn[T]
mu sync.Mutex
elements []T
closed bool
waitChan chan struct{} // if not nil, will be closed when new element is added
}
func NewQueue[T any](combineFn CombineFn[T]) *Queue[T] {
return &Queue[T]{
combineFn: combineFn,
}
}
func (q *Queue[T]) Len() int {
q.mu.Lock()
defer q.mu.Unlock()
return len(q.elements)
}
// Elements returns copy of the queue.
func (q *Queue[T]) Elements() []T {
q.mu.Lock()
defer q.mu.Unlock()
return slices.Clone(q.elements)
}
func (q *Queue[T]) Add(n T) {
q.mu.Lock()
defer q.mu.Unlock()
if q.closed {
panic("queue already closed")
}
for i, e := range q.elements {
if c, ok := q.combineFn(e, n); ok {
// No need to signal, since we are not adding new element.
q.elements[i] = c
return
}
}
q.elements = append(q.elements, n)
q.notifyWaiters()
}
// Must be called with lock held.
func (q *Queue[T]) notifyWaiters() {
if q.waitChan != nil {
// Wakes up all waiting goroutines (but also possibly zero, if they stopped waiting already).
close(q.waitChan)
q.waitChan = nil
}
}
var ErrClosed = errors.New("queue closed")
// Next returns the next element in the queue. If no element is available, Next will block until
// an element is added to the queue, or provided context is done.
// If the queue is closed, ErrClosed is returned.
func (q *Queue[T]) Next(ctx context.Context) (T, error) {
var zero T
q.mu.Lock()
unlockInDefer := true
defer func() {
if unlockInDefer {
q.mu.Unlock()
}
}()
for len(q.elements) == 0 {
if q.closed {
return zero, ErrClosed
}
// Wait for an element. Make sure there's a wait channel that we can use.
wch := q.waitChan
if wch == nil {
wch = make(chan struct{})
q.waitChan = wch
}
// Unlock before waiting
q.mu.Unlock()
select {
case <-ctx.Done():
unlockInDefer = false
return zero, ctx.Err()
case <-wch:
q.mu.Lock()
}
}
first := q.elements[0]
q.elements = q.elements[1:]
return first, nil
}
func (q *Queue[T]) Close() {
q.mu.Lock()
defer q.mu.Unlock()
q.closed = true
q.notifyWaiters()
}
+155
View File
@@ -0,0 +1,155 @@
package debouncer
import (
"context"
"errors"
"math/rand"
"sync"
"testing"
"time"
"github.com/stretchr/testify/require"
"go.uber.org/atomic"
"go.uber.org/goleak"
)
// This verifies that all goroutines spawned from tests are finished at the end of tests.
// Applies to all tests in the package.
func TestMain(m *testing.M) {
goleak.VerifyTestMain(m)
}
func TestQueueBasic(t *testing.T) {
q := NewQueue(func(a, b int) (c int, ok bool) {
return a + b, true
})
ctx, cancel := context.WithTimeout(context.Background(), 1*time.Millisecond)
defer cancel()
// Empty queue will time out.
require.Equal(t, context.DeadlineExceeded, nextErr(t, q, ctx))
q.Add(10)
require.Equal(t, 10, next(t, q))
require.Equal(t, 0, q.Len())
q.Add(20)
require.Equal(t, 20, next(t, q))
require.Equal(t, 0, q.Len())
q.Add(10)
require.Equal(t, 1, q.Len())
q.Add(20)
require.Equal(t, 1, q.Len())
require.Equal(t, 30, next(t, q))
require.Equal(t, 0, q.Len())
q.Add(100)
require.Equal(t, 1, q.Len())
q.Close()
require.Equal(t, 1, q.Len())
require.Equal(t, 100, next(t, q))
require.Equal(t, ErrClosed, nextErr(t, q, context.Background()))
require.Equal(t, 0, q.Len())
// We can call Next repeatedly, but will always get error.
require.Equal(t, ErrClosed, nextErr(t, q, context.Background()))
}
func TestQueueConcurrency(t *testing.T) {
q := NewQueue(func(a, b int64) (c int64, ok bool) {
// Combine the same numbers together.
if a == b {
return a + b, true
}
return 0, false
})
const numbers = 10000
const writeConcurrency = 50
const readConcurrency = 25
r := rand.New(rand.NewSource(time.Now().UnixNano()))
totalWrittenSum := atomic.NewInt64(0)
totalReadSum := atomic.NewInt64(0)
addCalls := atomic.NewInt64(0)
nextCalls := atomic.NewInt64(0)
// We will add some numbers to the queue.
writesWG := sync.WaitGroup{}
for i := 0; i < writeConcurrency; i++ {
writesWG.Add(1)
go func() {
defer writesWG.Done()
for j := 0; j < numbers; j++ {
v := r.Int63n(100) // Generate small number, so that we have a chance for combining some numbers.
q.Add(v)
addCalls.Inc()
totalWrittenSum.Add(v)
}
}()
}
readsWG := sync.WaitGroup{}
for i := 0; i < readConcurrency; i++ {
readsWG.Add(1)
go func() {
defer readsWG.Done()
for {
v, err := q.Next(context.Background())
if errors.Is(err, ErrClosed) {
return
}
require.NoError(t, err)
nextCalls.Inc()
totalReadSum.Add(v)
}
}()
}
writesWG.Wait()
// Close queue after sending all numbers. This signals readers that they can stop.
q.Close()
// Wait until all readers finish too.
readsWG.Wait()
// Verify that all numbers were sent, combined and received.
require.Equal(t, int64(writeConcurrency*numbers), addCalls.Load())
require.Equal(t, totalWrittenSum.Load(), totalReadSum.Load())
require.LessOrEqual(t, nextCalls.Load(), addCalls.Load())
t.Log("add calls:", addCalls.Load(), "next calls:", nextCalls.Load(), "total written sum:", totalWrittenSum.Load(), "total read sum:", totalReadSum.Load())
}
func TestQueueCloseUnblocksReaders(t *testing.T) {
q := NewQueue(func(a, b int) (c int, ok bool) {
return a + b, true
})
wg := sync.WaitGroup{}
wg.Add(1)
go func() {
defer wg.Done()
time.Sleep(50 * time.Millisecond)
q.Close()
}()
_, err := q.Next(context.Background())
require.ErrorIs(t, err, ErrClosed)
wg.Wait()
}
func next[T any](t *testing.T, q *Queue[T]) T {
v, err := q.Next(context.Background())
require.NoError(t, err)
return v
}
func nextErr[T any](t *testing.T, q *Queue[T], ctx context.Context) error {
_, err := q.Next(ctx)
require.Error(t, err)
return err
}