diff --git a/pkg/storage/unified/resource/search_queue.go b/pkg/storage/unified/resource/search_queue.go index 9c71d913222..23310caf9df 100644 --- a/pkg/storage/unified/resource/search_queue.go +++ b/pkg/storage/unified/resource/search_queue.go @@ -12,7 +12,6 @@ import ( // It is responsible for ingesting events for a single index // It will batch events and send them to the index in a single bulk request type indexQueueProcessor struct { - index ResourceIndex nsr NamespacedResource queue chan *WrittenEvent batchSize int @@ -20,8 +19,11 @@ type indexQueueProcessor struct { resChan chan *IndexEvent // Channel to send results to the caller - mu sync.Mutex - running bool + indexMu sync.Mutex + index ResourceIndex + + runningMu sync.Mutex + running bool } type IndexEvent struct { @@ -34,8 +36,8 @@ type IndexEvent struct { } func (b *indexQueueProcessor) updateIndex(newIndex ResourceIndex) { - b.mu.Lock() - defer b.mu.Unlock() + b.indexMu.Lock() + defer b.indexMu.Unlock() b.index = newIndex } @@ -57,8 +59,8 @@ func (b *indexQueueProcessor) Add(evt *WrittenEvent) { b.queue <- evt // Start the processor if it's not already running - b.mu.Lock() - defer b.mu.Unlock() + b.runningMu.Lock() + defer b.runningMu.Unlock() if !b.running { b.running = true go b.runProcessor() @@ -68,9 +70,9 @@ func (b *indexQueueProcessor) Add(evt *WrittenEvent) { // runProcessor is the task processing the queue of written events func (b *indexQueueProcessor) runProcessor() { defer func() { - b.mu.Lock() + b.runningMu.Lock() b.running = false - b.mu.Unlock() + b.runningMu.Unlock() }() for { @@ -122,8 +124,9 @@ func (b *indexQueueProcessor) process(batch []*WrittenEvent) { } else { item.Action = ActionIndex ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cancel() doc, err := b.builder.BuildDocument(ctx, evt.Key, evt.ResourceVersion, evt.Value) + cancel() + if err != nil { result.Err = err } else { @@ -134,10 +137,11 @@ func (b *indexQueueProcessor) process(batch []*WrittenEvent) { req.Items = append(req.Items, item) } - b.mu.Lock() - err := b.index.BulkIndex(req) - b.mu.Unlock() + b.indexMu.Lock() + idx := b.index + b.indexMu.Unlock() + err := idx.BulkIndex(req) if err != nil { for _, r := range resp { r.Err = err