From ed6bfc0b67cba2264d4c7698aa55f3f1896e171d Mon Sep 17 00:00:00 2001 From: Blaine Motsinger Date: Mon, 28 Sep 2026 12:16:32 -0400 Subject: [PATCH] GH-31: fix index race: async open, bounded tool waits --- main.go | 15 +++-- search/index.go | 154 +++++++++++++++++++++++++++++++++++++------- search/race_test.go | 125 +++++++++++++++++++++++++++++++++++ 3 files changed, 267 insertions(+), 27 deletions(-) create mode 100644 search/race_test.go diff --git a/main.go b/main.go index 8dc8eff..c6b4af9 100644 --- a/main.go +++ b/main.go @@ -94,15 +94,20 @@ func main() { log.Fatalf("determine index location: %v", err) } - idx, err := search.NewIndex(indexBasePath, resolver.KnownNames()) - if err != nil { - log.Fatalf("init search index: %v", err) - } + idx := search.NewLazyIndex(indexBasePath, resolver.KnownNames()) defer idx.Close() + // Open the index in the background: a second instance whose index is + // locked by another knowledge-mcp (single-writer) must not delay the + // MCP protocol. Tool calls degrade to a clear error until the lock + // frees, after which the pending open completes on its own. + var afterOpen func() if !*noIndexOnStartup { - go idx.IndexAll(resolver) + afterOpen = func() { + idx.IndexAll(resolver) + } } + idx.OpenBackground(afterOpen) // Create MCP server s := server.NewMCPServer( diff --git a/search/index.go b/search/index.go index 68fc373..2691b40 100644 --- a/search/index.go +++ b/search/index.go @@ -2,6 +2,7 @@ package search import ( "encoding/json" + "errors" "fmt" "log" "os" @@ -10,6 +11,7 @@ import ( "sort" "strings" "sync" + "time" "github.com/blevesearch/bleve/v2" "github.com/blevesearch/bleve/v2/mapping" @@ -56,22 +58,115 @@ type SearchResult struct { // Index wraps a bleve index for knowledge search. type Index struct { indexPath string + names []string index bleve.Index ready chan struct{} // closed when indexing completes or stale data exists readyOnce sync.Once // guards close of ready to prevent double-close panic + + openMu sync.Mutex + openStarted bool // guards against double opens + openDone chan struct{} // closed when the open attempt resolves + openErr error // set when the open attempt fails; read after openDone closes } -// NewIndex opens or creates a bleve index at the given path. indexNames is the -// set of project names the index will hold; a change from the recorded set -// triggers a rebuild. +// openWaitTimeout bounds how long Query/Add wait on a pending index open. +// The only expected long-open case is another instance holding the +// single-writer bbolt lock; everything else opens in milliseconds. +const openWaitTimeout = 2 * time.Second + +// ErrIndexLocked is returned when the index is not open within the bounded +// wait — most commonly because another instance holds the single-writer lock. +var ErrIndexLocked = errors.New("search index is not ready; another knowledge-mcp instance may hold the single-writer index lock, retry in a moment") + +// NewIndex opens or creates a bleve index at the given path synchronously. +// indexNames is the set of project names the index will hold; a change from +// the recorded set triggers a rebuild. func NewIndex(indexPath string, indexNames []string) (*Index, error) { - if err := os.MkdirAll(filepath.Dir(indexPath), 0755); err != nil { - return nil, fmt.Errorf("create index dir: %w", err) + i := NewLazyIndex(indexPath, indexNames) + i.openStarted = true + if err := i.open(); err != nil { + return nil, err } + close(i.openDone) + return i, nil +} + +// NewLazyIndex constructs an Index without touching disk. The open is +// deferred until OpenBackground runs it. +func NewLazyIndex(indexPath string, indexNames []string) *Index { + return &Index{ + indexPath: indexPath, + names: indexNames, + ready: make(chan struct{}), + openDone: make(chan struct{}), + } +} - if metaStale(indexPath, indexNames) { - if err := DeleteIndex(indexPath); err != nil { - return nil, fmt.Errorf("remove stale index: %w", err) +// OpenBackground starts the index open on a background goroutine exactly +// once. afterOpen runs after a successful open. The open may block for a +// long time if another instance holds the single-writer lock; that must not +// delay the caller (the MCP server) from serving its protocol. +func (i *Index) OpenBackground(afterOpen func()) { + i.openMu.Lock() + if i.openStarted { + i.openMu.Unlock() + return + } + i.openStarted = true + i.openMu.Unlock() + + go func() { + err := i.open() + i.openMu.Lock() + i.openErr = err + i.openMu.Unlock() + close(i.openDone) + if err == nil && afterOpen != nil { + afterOpen() + } + }() +} + +// WaitOpen waits up to d for the open attempt to resolve. It returns nil +// once the index is open, the open error if the attempt failed, or +// ErrIndexLocked if the attempt is still pending after d. +func (i *Index) WaitOpen(d time.Duration) error { + waitCh := func() error { + if i.openErr != nil { + return fmt.Errorf("search index unavailable: %w", i.openErr) + } + return nil + } + + timer := time.NewTimer(d) + defer timer.Stop() + select { + case <-i.openDone: + i.openMu.Lock() + defer i.openMu.Unlock() + return waitCh() + case <-timer.C: + select { + case <-i.openDone: + i.openMu.Lock() + defer i.openMu.Unlock() + return waitCh() + default: + return ErrIndexLocked + } + } +} + +// open performs the actual on-disk open or create, exactly as NewIndex +// historically did. +func (i *Index) open() error { + if err := os.MkdirAll(filepath.Dir(i.indexPath), 0755); err != nil { + return fmt.Errorf("create index dir: %w", err) + } + + if metaStale(i.indexPath, i.names) { + if err := DeleteIndex(i.indexPath); err != nil { + return fmt.Errorf("remove stale index: %w", err) } } @@ -80,36 +175,38 @@ func NewIndex(indexPath string, indexNames []string) (*Index, error) { var idx bleve.Index var err error - if _, statErr := os.Stat(indexPath); os.IsNotExist(statErr) { - idx, err = bleve.New(indexPath, mapping) + if _, statErr := os.Stat(i.indexPath); os.IsNotExist(statErr) { + idx, err = bleve.New(i.indexPath, mapping) if err != nil { - return nil, fmt.Errorf("create index: %w", err) + return fmt.Errorf("create index: %w", err) } } else { - idx, err = openOrRecover(indexPath, mapping) + idx, err = openOrRecover(i.indexPath, mapping) if err != nil { - return nil, err + return err } } - meta := indexMeta{Version: metaFileVersion, Names: indexNames} + meta := indexMeta{Version: metaFileVersion, Names: i.names} data, marshalErr := json.Marshal(meta) if marshalErr == nil { - marshalErr = os.WriteFile(metaFilePath(indexPath), data, 0644) + marshalErr = os.WriteFile(metaFilePath(i.indexPath), data, 0644) } if marshalErr != nil { idx.Close() - return nil, fmt.Errorf("write index meta: %w", marshalErr) + return fmt.Errorf("write index meta: %w", marshalErr) } - ki := &Index{indexPath: indexPath, index: idx, ready: make(chan struct{})} + i.openMu.Lock() + i.index = idx + i.openMu.Unlock() // If index has data, close ready immediately (stale data available) - if ki.indexHasData() { - ki.CloseReady() + if i.indexHasData() { + i.CloseReady() } - return ki, nil + return nil } func metaFilePath(indexPath string) string { @@ -202,12 +299,19 @@ func indexMapping() mapping.IndexMapping { // Add indexes a document with the given (scoped) key. func (i *Index) Add(id string, doc SearchDocument) error { + if err := i.WaitOpen(openWaitTimeout); err != nil { + return err + } return i.index.Index(id, doc) } // Query searches the index, scoped to one project, with full-text // search and optional filters. func (i *Index) Query(project, q, category string, limit int) ([]SearchResult, error) { + if err := i.WaitOpen(openWaitTimeout); err != nil { + return nil, err + } + // Block until indexing completes if no stale data available <-i.ready @@ -282,9 +386,15 @@ func (i *Index) Query(project, q, category string, limit int) ([]SearchResult, e return results, nil } -// Close closes the bleve index. +// Close closes the bleve index, if it was ever opened. func (i *Index) Close() error { - return i.index.Close() + i.openMu.Lock() + idx := i.index + i.openMu.Unlock() + if idx == nil { + return nil + } + return idx.Close() } // CloseReady closes the ready channel exactly once, preventing double-close panics. diff --git a/search/race_test.go b/search/race_test.go new file mode 100644 index 0000000..2c2f1d5 --- /dev/null +++ b/search/race_test.go @@ -0,0 +1,125 @@ +package search + +import ( + "path/filepath" + "testing" + "time" +) + +// TestOpenBackgroundDoesNotBlockWhileLocked covers the startup race: a second +// server instance must not wedge the MCP handshake when the first instance +// holds the single-writer index lock (bbolt flocks block indefinitely). +func TestOpenBackgroundDoesNotBlockWhileLocked(t *testing.T) { + dir := t.TempDir() + indexPath := filepath.Join(dir, "shared.index") + + // Instance A creates the index and holds the bbolt lock. + a := newIndexReady(t, indexPath, []string{"test"}) + if err := a.Add("test/conv-001", SearchDocument{Summary: "persisted doc", Project: "test"}); err != nil { + t.Fatalf("Add() to instance A: %v", err) + } + + // Instance B attempts an async open; it must not block this goroutine. + b := NewLazyIndex(indexPath, []string{"test"}) + b.OpenBackground(nil) + defer b.Close() + + if err := b.WaitOpen(250 * time.Millisecond); err == nil { + t.Fatal("WaitOpen() resolved while the index lock was still held, want unresolved") + } + + // Once A releases the lock, B's pending open must complete on its own. + if err := a.Close(); err != nil { + t.Fatalf("closing instance A: %v", err) + } + if err := b.WaitOpen(10 * time.Second); err != nil { + t.Fatalf("WaitOpen() after lock release: %v", err) + } + + results, err := b.Query("test", "persisted", "", 10) + if err != nil { + t.Fatalf("Query() after recovery: %v", err) + } + if len(results) != 1 { + t.Fatalf("Query() returned %d results after recovery, want 1", len(results)) + } +} + +// TestQueryDoesNotHangWhileLocked covers tool calls: a query against a +// locked index must return an error promptly instead of blocking forever. +func TestQueryDoesNotHangWhileLocked(t *testing.T) { + dir := t.TempDir() + indexPath := filepath.Join(dir, "shared.index") + + a := newIndexReady(t, indexPath, []string{"test"}) + if err := a.Add("test/conv-001", SearchDocument{Summary: "locked doc", Project: "test"}); err != nil { + t.Fatalf("Add() to instance A: %v", err) + } + + b := NewLazyIndex(indexPath, []string{"test"}) + b.OpenBackground(nil) + defer b.Close() + + done := make(chan error, 1) + go func() { + _, err := b.Query("test", "locked", "", 10) + done <- err + }() + + select { + case err := <-done: + if err == nil { + t.Fatal("Query() succeeded against a locked index, want error") + } + case <-time.After(5 * time.Second): + t.Fatal("Query() hung while the index was locked, want prompt error") + } + + // After the lock is released the same query must succeed, proving the + // instance recovers without a restart. + if err := a.Close(); err != nil { + t.Fatalf("closing instance A: %v", err) + } + if err := b.WaitOpen(10 * time.Second); err != nil { + t.Fatalf("WaitOpen() after lock release: %v", err) + } + results, err := b.Query("test", "locked", "", 10) + if err != nil { + t.Fatalf("Query() after recovery: %v", err) + } + if len(results) != 1 { + t.Fatalf("Query() returned %d results after recovery, want 1", len(results)) + } +} + +// TestAddDoesNotPanicWhileLocked covers writes: Add must surface the open +// state as an error, never dereference a nil index. +func TestAddDoesNotPanicWhileLocked(t *testing.T) { + dir := t.TempDir() + indexPath := filepath.Join(dir, "shared.index") + + a := newIndexReady(t, indexPath, []string{"test"}) + + b := NewLazyIndex(indexPath, []string{"test"}) + b.OpenBackground(nil) + defer b.Close() + + errCh := make(chan error, 1) + go func() { + errCh <- b.Add("test/conv-002", SearchDocument{Summary: "new doc", Project: "test"}) + }() + + select { + case err := <-errCh: + if err == nil { + t.Fatal("Add() succeeded against a locked index, want error") + } + case <-time.After(5 * time.Second): + t.Fatal("Add() hung while the index was locked, want prompt error") + } + + if err := a.Close(); err != nil { + t.Fatalf("closing instance A: %v", err) + } + _ = b.WaitOpen(10 * time.Second) +}