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
15 changes: 10 additions & 5 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
154 changes: 132 additions & 22 deletions search/index.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package search

import (
"encoding/json"
"errors"
"fmt"
"log"
"os"
Expand All @@ -10,6 +11,7 @@ import (
"sort"
"strings"
"sync"
"time"

"github.com/blevesearch/bleve/v2"
"github.com/blevesearch/bleve/v2/mapping"
Expand Down Expand Up @@ -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)
}
}

Expand All @@ -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 {
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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.
Expand Down
125 changes: 125 additions & 0 deletions search/race_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
Loading