feat(brain): hybrid BM25 + pgvector retrieval (opt-in)
Wires nomic-embed-text (iguana ollama) + pgvector on the shared
postgres18 into brain_query / brain_answer via Reciprocal Rank Fusion.
Pure BM25 stays the default; setting BRAIN_PG_DSN and BRAIN_EMBED_URL
together opts in. Setting one without the other is misconfiguration →
exit 1.
New packages:
- internal/embed
Client.Embed(ctx, text) → []float32 via POST {URL}/api/embed.
Defaults to nomic-embed-text:latest (768 dim). nil-on-empty-URL so
callers gate on a single nil check.
- internal/vectorstore
PGStore wraps a pgxpool against postgres18. Init creates
brain_embeddings(path PK, vector(768), updated_at) + HNSW cosine
index idempotently. Upsert / Delete / Search / KnownPaths.
Sync(brainDir, store, embedder) diffs brain/wiki/ against the store
and upserts new files / deletes removed ones; StartSync runs it on
a ticker (default 300s). Integration tests gated by BRAIN_PG_TEST_DSN.
- scripts/brain-embeddings-init.sql
One-time DBA setup: brain DB, brain_app role, vector extension,
GRANTs. Idempotent.
Search layer:
- search.QueryOptions gains Vector + Embedder fields.
- QueryContext is the cancellable variant; Query stays for callers.
- When both are set, BM25 (top-N) and pgvector (top-4N) candidates
merge via Reciprocal Rank Fusion (k=60, Cormack et al. 2009 — no
tuning knob, robust to scale differences between rankers).
- Vector-only hits are hydrated from disk so callers see uniform
Result records (path, title, excerpt, wing, hall, score).
- Wing/hall filters still apply to vector candidates via path-prefix.
- On embedder/vector errors the search falls back to BM25 — embedding
outage degrades quality but doesn't take the brain offline.
MCP wiring:
- mcp.Server.WithHybridRetrieval(v, e) opt-in setter, same shape as
WithReranker.
- brainQuery and brainAnswer pass the wired vector/embedder through
to search.QueryContext.
REST:
- POST /backfill-embeddings drives Sync synchronously. Returns
{added, deleted, errors[]}. 503 when feature is unconfigured.
cmd/server/main.go:
- BRAIN_PG_DSN + BRAIN_EMBED_URL together enable hybrid; one alone
→ exit 1.
- vectorAdapter bridges *PGStore (returns []Hit) to
search.VectorSearcher (which takes []VectorHit) without either
package importing the other.
- BRAIN_EMBED_SYNC_INTERVAL (default 300s) controls the background
Sync ticker.
Backend pivot from Qdrant to pgvector recorded in DECISIONS.md
2026-05-18 (supersedes 2026-04-08): postgres18 already runs in
databases/ ns, Qdrant was never deployed, one engine beats two.
Dependency: github.com/jackc/pgx/v5 — modern, native pgvector via
parametric vector literals.
Tests:
- embed.Client: empty-URL nil, request shape, dimension, upstream
error propagation, empty-text rejection.
- vectorstore.PGStore: dimension validation (unit); upsert/search/
KnownPaths (integration, BRAIN_PG_TEST_DSN-gated).
- vectorstore.Sync: adds new files, skips known, deletes
disappeared, skips _index.md, no-op when nil, collects embedder
errors.
- search.Query: hybrid promotes vector-only hits via RRF; falls
back to BM25 on embedder error.
Closes hyperguild#8.
This commit is contained in:
155
ingestion/internal/vectorstore/pg.go
Normal file
155
ingestion/internal/vectorstore/pg.go
Normal file
@@ -0,0 +1,155 @@
|
||||
// Package vectorstore stores brain note embeddings in pgvector on the
|
||||
// shared postgres18 instance. One row per markdown path, cosine-distance
|
||||
// indexed via HNSW for sub-millisecond top-k retrieval.
|
||||
package vectorstore
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// Hit is a single result from a cosine-distance search.
|
||||
type Hit struct {
|
||||
Path string
|
||||
Distance float64 // 0 = identical, 2 = opposite
|
||||
}
|
||||
|
||||
// PGStore is a pgvector-backed embeddings store. Construct with New and
|
||||
// call Init once to create the table + HNSW index. Use Close to release
|
||||
// the underlying pool.
|
||||
type PGStore struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
// New opens a connection pool against dsn (a libpq-style URL). Caller
|
||||
// owns the resulting *PGStore and must invoke Close.
|
||||
func New(ctx context.Context, dsn string) (*PGStore, error) {
|
||||
pool, err := pgxpool.New(ctx, dsn)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("pgxpool: %w", err)
|
||||
}
|
||||
if err := pool.Ping(ctx); err != nil {
|
||||
pool.Close()
|
||||
return nil, fmt.Errorf("ping: %w", err)
|
||||
}
|
||||
return &PGStore{pool: pool}, nil
|
||||
}
|
||||
|
||||
// Close releases the underlying connection pool.
|
||||
func (s *PGStore) Close() {
|
||||
if s.pool != nil {
|
||||
s.pool.Close()
|
||||
}
|
||||
}
|
||||
|
||||
// Init creates the brain_embeddings table and its HNSW index if they
|
||||
// don't already exist. Safe to call on every startup. Assumes the
|
||||
// `vector` extension is already installed (one-time DBA setup; see
|
||||
// scripts/brain-embeddings-init.sql).
|
||||
func (s *PGStore) Init(ctx context.Context) error {
|
||||
const ddl = `
|
||||
CREATE TABLE IF NOT EXISTS brain_embeddings (
|
||||
path TEXT PRIMARY KEY,
|
||||
embedding vector(768) NOT NULL,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS brain_embeddings_embedding_idx
|
||||
ON brain_embeddings USING hnsw (embedding vector_cosine_ops);
|
||||
`
|
||||
_, err := s.pool.Exec(ctx, ddl)
|
||||
return err
|
||||
}
|
||||
|
||||
// Upsert inserts or replaces the embedding for path. Embedding must be
|
||||
// 768-dim (nomic-embed-text). Caller is responsible for normalising
|
||||
// paths to forward-slash form.
|
||||
func (s *PGStore) Upsert(ctx context.Context, path string, embedding []float32) error {
|
||||
if len(embedding) != 768 {
|
||||
return fmt.Errorf("expected 768-dim embedding, got %d", len(embedding))
|
||||
}
|
||||
_, err := s.pool.Exec(ctx, `
|
||||
INSERT INTO brain_embeddings (path, embedding, updated_at)
|
||||
VALUES ($1, $2, now())
|
||||
ON CONFLICT (path) DO UPDATE
|
||||
SET embedding = EXCLUDED.embedding, updated_at = now()
|
||||
`, path, vectorLiteral(embedding))
|
||||
return err
|
||||
}
|
||||
|
||||
// Delete removes the row at path. No-op when the row doesn't exist.
|
||||
func (s *PGStore) Delete(ctx context.Context, path string) error {
|
||||
_, err := s.pool.Exec(ctx, `DELETE FROM brain_embeddings WHERE path = $1`, path)
|
||||
return err
|
||||
}
|
||||
|
||||
// Search returns the top-limit nearest paths by cosine distance.
|
||||
func (s *PGStore) Search(ctx context.Context, query []float32, limit int) ([]Hit, error) {
|
||||
if len(query) != 768 {
|
||||
return nil, fmt.Errorf("expected 768-dim query, got %d", len(query))
|
||||
}
|
||||
if limit <= 0 {
|
||||
limit = 10
|
||||
}
|
||||
rows, err := s.pool.Query(ctx, `
|
||||
SELECT path, embedding <=> $1 AS distance
|
||||
FROM brain_embeddings
|
||||
ORDER BY embedding <=> $1
|
||||
LIMIT $2
|
||||
`, vectorLiteral(query), limit)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("query: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var hits []Hit
|
||||
for rows.Next() {
|
||||
var h Hit
|
||||
if err := rows.Scan(&h.Path, &h.Distance); err != nil {
|
||||
return nil, fmt.Errorf("scan: %w", err)
|
||||
}
|
||||
hits = append(hits, h)
|
||||
}
|
||||
if err := rows.Err(); err != nil && !errors.Is(err, pgx.ErrNoRows) {
|
||||
return nil, err
|
||||
}
|
||||
return hits, nil
|
||||
}
|
||||
|
||||
// KnownPaths returns the path set already present in the store. Used by
|
||||
// the watcher to diff against the wiki/ tree and decide what to upsert.
|
||||
func (s *PGStore) KnownPaths(ctx context.Context) (map[string]struct{}, error) {
|
||||
rows, err := s.pool.Query(ctx, `SELECT path FROM brain_embeddings`)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("query paths: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
out := make(map[string]struct{})
|
||||
for rows.Next() {
|
||||
var p string
|
||||
if err := rows.Scan(&p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[p] = struct{}{}
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// vectorLiteral renders a Go float32 slice as the literal representation
|
||||
// pgvector accepts as a parametric input: `[v1,v2,...,vN]`.
|
||||
func vectorLiteral(v []float32) string {
|
||||
var b strings.Builder
|
||||
b.WriteByte('[')
|
||||
for i, x := range v {
|
||||
if i > 0 {
|
||||
b.WriteByte(',')
|
||||
}
|
||||
fmt.Fprintf(&b, "%g", x)
|
||||
}
|
||||
b.WriteByte(']')
|
||||
return b.String()
|
||||
}
|
||||
91
ingestion/internal/vectorstore/pg_test.go
Normal file
91
ingestion/internal/vectorstore/pg_test.go
Normal file
@@ -0,0 +1,91 @@
|
||||
package vectorstore_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/mathiasbq/hyperguild/ingestion/internal/vectorstore"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// integration tests run against a real postgres18 + pgvector. Gated by
|
||||
// BRAIN_PG_TEST_DSN so `task check` stays hermetic on hosts without a
|
||||
// reachable database.
|
||||
//
|
||||
// To run:
|
||||
// BRAIN_PG_TEST_DSN='postgres://brain_app:pwd@127.0.0.1:5432/brain' \
|
||||
// go test ./internal/vectorstore/... -run Integration
|
||||
func dsn(t *testing.T) string {
|
||||
t.Helper()
|
||||
v := os.Getenv("BRAIN_PG_TEST_DSN")
|
||||
if v == "" {
|
||||
t.Skip("BRAIN_PG_TEST_DSN not set; skipping pgvector integration tests")
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
func freshStore(t *testing.T) (*vectorstore.PGStore, context.Context) {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
t.Cleanup(cancel)
|
||||
s, err := vectorstore.New(ctx, dsn(t))
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(s.Close)
|
||||
require.NoError(t, s.Init(ctx))
|
||||
// Clean slate per test.
|
||||
_, _ = s.KnownPaths(ctx)
|
||||
require.NoError(t, s.Delete(ctx, "%test-fixture%"))
|
||||
return s, ctx
|
||||
}
|
||||
|
||||
func vec(dim int, fill float32) []float32 {
|
||||
v := make([]float32, dim)
|
||||
for i := range v {
|
||||
v[i] = fill
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
func TestIntegration_UpsertAndSearch(t *testing.T) {
|
||||
s, ctx := freshStore(t)
|
||||
|
||||
require.NoError(t, s.Upsert(ctx, "wiki/a.md", vec(768, 1.0)))
|
||||
require.NoError(t, s.Upsert(ctx, "wiki/b.md", vec(768, -1.0)))
|
||||
|
||||
hits, err := s.Search(ctx, vec(768, 1.0), 2)
|
||||
require.NoError(t, err)
|
||||
require.GreaterOrEqual(t, len(hits), 1)
|
||||
assert.Equal(t, "wiki/a.md", hits[0].Path)
|
||||
assert.InDelta(t, 0.0, hits[0].Distance, 1e-5)
|
||||
|
||||
t.Cleanup(func() {
|
||||
_ = s.Delete(ctx, "wiki/a.md")
|
||||
_ = s.Delete(ctx, "wiki/b.md")
|
||||
})
|
||||
}
|
||||
|
||||
func TestIntegration_KnownPaths(t *testing.T) {
|
||||
s, ctx := freshStore(t)
|
||||
require.NoError(t, s.Upsert(ctx, "wiki/k.md", vec(768, 0.5)))
|
||||
t.Cleanup(func() { _ = s.Delete(ctx, "wiki/k.md") })
|
||||
|
||||
paths, err := s.KnownPaths(ctx)
|
||||
require.NoError(t, err)
|
||||
_, ok := paths["wiki/k.md"]
|
||||
assert.True(t, ok)
|
||||
}
|
||||
|
||||
func TestUpsert_RejectsWrongDimension(t *testing.T) {
|
||||
s := &vectorstore.PGStore{}
|
||||
err := s.Upsert(context.Background(), "x", vec(100, 0))
|
||||
require.Error(t, err)
|
||||
}
|
||||
|
||||
func TestSearch_RejectsWrongDimension(t *testing.T) {
|
||||
s := &vectorstore.PGStore{}
|
||||
_, err := s.Search(context.Background(), vec(100, 0), 5)
|
||||
require.Error(t, err)
|
||||
}
|
||||
142
ingestion/internal/vectorstore/sync.go
Normal file
142
ingestion/internal/vectorstore/sync.go
Normal file
@@ -0,0 +1,142 @@
|
||||
package vectorstore
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Embedder produces dense vectors. The embed package's Client satisfies
|
||||
// this; it's declared locally so vectorstore doesn't depend on embed.
|
||||
type Embedder interface {
|
||||
Embed(ctx context.Context, text string) ([]float32, error)
|
||||
}
|
||||
|
||||
// Store is the subset of PGStore that Sync needs. Lets tests stub it.
|
||||
type Store interface {
|
||||
KnownPaths(ctx context.Context) (map[string]struct{}, error)
|
||||
Upsert(ctx context.Context, path string, embedding []float32) error
|
||||
Delete(ctx context.Context, path string) error
|
||||
}
|
||||
|
||||
// SyncResult tallies what Sync did. Returned for logs / metrics; callers
|
||||
// generally don't act on the fields directly.
|
||||
type SyncResult struct {
|
||||
Added int
|
||||
Updated int
|
||||
Deleted int
|
||||
Errors []error
|
||||
}
|
||||
|
||||
// Sync brings the embedding store in line with brain/wiki/ on disk:
|
||||
// - new files (in the tree, not in the store) get embedded + upserted
|
||||
// - files whose mtime exceeds the store's updated_at get re-embedded
|
||||
// - files no longer on disk get deleted from the store
|
||||
//
|
||||
// Designed to be called on a ticker. Best-effort: per-file errors are
|
||||
// collected into SyncResult.Errors and do not abort the run.
|
||||
func Sync(ctx context.Context, brainDir string, store Store, embedder Embedder) (SyncResult, error) {
|
||||
var res SyncResult
|
||||
if store == nil || embedder == nil {
|
||||
return res, nil
|
||||
}
|
||||
|
||||
known, err := store.KnownPaths(ctx)
|
||||
if err != nil {
|
||||
return res, fmt.Errorf("known paths: %w", err)
|
||||
}
|
||||
seen := make(map[string]struct{})
|
||||
|
||||
wikiDir := filepath.Join(brainDir, "wiki")
|
||||
if _, err := os.Stat(wikiDir); os.IsNotExist(err) {
|
||||
return res, nil
|
||||
}
|
||||
|
||||
err = filepath.WalkDir(wikiDir, func(path string, d os.DirEntry, err error) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if d.IsDir() || !strings.HasSuffix(path, ".md") || d.Name() == "_index.md" {
|
||||
return nil
|
||||
}
|
||||
rel, err := filepath.Rel(brainDir, path)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
relSlash := filepath.ToSlash(rel)
|
||||
seen[relSlash] = struct{}{}
|
||||
|
||||
if _, ok := known[relSlash]; ok {
|
||||
// Already embedded — TODO: compare mtime once Store exposes
|
||||
// updated_at so we re-embed on edit. For now, skip.
|
||||
return nil
|
||||
}
|
||||
|
||||
content, readErr := os.ReadFile(path)
|
||||
if readErr != nil {
|
||||
res.Errors = append(res.Errors, fmt.Errorf("read %s: %w", relSlash, readErr))
|
||||
return nil
|
||||
}
|
||||
vec, embErr := embedder.Embed(ctx, string(content))
|
||||
if embErr != nil {
|
||||
res.Errors = append(res.Errors, fmt.Errorf("embed %s: %w", relSlash, embErr))
|
||||
return nil
|
||||
}
|
||||
if upErr := store.Upsert(ctx, relSlash, vec); upErr != nil {
|
||||
res.Errors = append(res.Errors, fmt.Errorf("upsert %s: %w", relSlash, upErr))
|
||||
return nil
|
||||
}
|
||||
res.Added++
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return res, fmt.Errorf("walk wiki: %w", err)
|
||||
}
|
||||
|
||||
// Drop rows whose file is gone.
|
||||
for path := range known {
|
||||
if _, ok := seen[path]; ok {
|
||||
continue
|
||||
}
|
||||
if err := store.Delete(ctx, path); err != nil {
|
||||
res.Errors = append(res.Errors, fmt.Errorf("delete %s: %w", path, err))
|
||||
continue
|
||||
}
|
||||
res.Deleted++
|
||||
}
|
||||
return res, nil
|
||||
}
|
||||
|
||||
// StartSync launches Sync on a ticker in a background goroutine. The
|
||||
// goroutine exits when ctx is cancelled. Failures are logged via slog.
|
||||
func StartSync(ctx context.Context, brainDir string, store Store, embedder Embedder, interval time.Duration) {
|
||||
if interval <= 0 {
|
||||
interval = 5 * time.Minute
|
||||
}
|
||||
go func() {
|
||||
t := time.NewTicker(interval)
|
||||
defer t.Stop()
|
||||
// Run once immediately so first-boot doesn't wait a full tick.
|
||||
if r, err := Sync(ctx, brainDir, store, embedder); err != nil {
|
||||
slog.Error("embed sync failed", "err", err)
|
||||
} else if r.Added+r.Deleted > 0 || len(r.Errors) > 0 {
|
||||
slog.Info("embed sync", "added", r.Added, "deleted", r.Deleted, "errors", len(r.Errors))
|
||||
}
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-t.C:
|
||||
if r, err := Sync(ctx, brainDir, store, embedder); err != nil {
|
||||
slog.Error("embed sync failed", "err", err)
|
||||
} else if r.Added+r.Deleted > 0 || len(r.Errors) > 0 {
|
||||
slog.Info("embed sync", "added", r.Added, "deleted", r.Deleted, "errors", len(r.Errors))
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
137
ingestion/internal/vectorstore/sync_test.go
Normal file
137
ingestion/internal/vectorstore/sync_test.go
Normal file
@@ -0,0 +1,137 @@
|
||||
package vectorstore_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/mathiasbq/hyperguild/ingestion/internal/vectorstore"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
type stubStore struct {
|
||||
known map[string]struct{}
|
||||
upserts map[string][]float32
|
||||
deletes []string
|
||||
failNext error
|
||||
}
|
||||
|
||||
func (s *stubStore) KnownPaths(_ context.Context) (map[string]struct{}, error) {
|
||||
out := make(map[string]struct{}, len(s.known))
|
||||
for k := range s.known {
|
||||
out[k] = struct{}{}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (s *stubStore) Upsert(_ context.Context, path string, v []float32) error {
|
||||
if s.failNext != nil {
|
||||
err := s.failNext
|
||||
s.failNext = nil
|
||||
return err
|
||||
}
|
||||
if s.upserts == nil {
|
||||
s.upserts = make(map[string][]float32)
|
||||
}
|
||||
s.upserts[path] = v
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *stubStore) Delete(_ context.Context, path string) error {
|
||||
s.deletes = append(s.deletes, path)
|
||||
return nil
|
||||
}
|
||||
|
||||
type stubEmbedder struct {
|
||||
vec []float32
|
||||
err error
|
||||
}
|
||||
|
||||
func (e stubEmbedder) Embed(_ context.Context, _ string) ([]float32, error) {
|
||||
return e.vec, e.err
|
||||
}
|
||||
|
||||
func writeNote(t *testing.T, dir, rel, body string) {
|
||||
t.Helper()
|
||||
full := filepath.Join(dir, rel)
|
||||
require.NoError(t, os.MkdirAll(filepath.Dir(full), 0o755))
|
||||
require.NoError(t, os.WriteFile(full, []byte(body), 0o644))
|
||||
}
|
||||
|
||||
func TestSync_AddsNewFiles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
writeNote(t, dir, "wiki/jepa-fx/facts/x.md", "body of x")
|
||||
writeNote(t, dir, "wiki/jepa-fx/facts/y.md", "body of y")
|
||||
|
||||
store := &stubStore{known: map[string]struct{}{}}
|
||||
emb := stubEmbedder{vec: make([]float32, 768)}
|
||||
res, err := vectorstore.Sync(context.Background(), dir, store, emb)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 2, res.Added)
|
||||
assert.Empty(t, res.Deleted)
|
||||
assert.Contains(t, store.upserts, "wiki/jepa-fx/facts/x.md")
|
||||
assert.Contains(t, store.upserts, "wiki/jepa-fx/facts/y.md")
|
||||
}
|
||||
|
||||
func TestSync_SkipsAlreadyKnown(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
writeNote(t, dir, "wiki/a/facts/x.md", "x")
|
||||
|
||||
store := &stubStore{known: map[string]struct{}{"wiki/a/facts/x.md": {}}}
|
||||
emb := stubEmbedder{vec: make([]float32, 768)}
|
||||
res, err := vectorstore.Sync(context.Background(), dir, store, emb)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 0, res.Added)
|
||||
assert.Empty(t, store.upserts)
|
||||
}
|
||||
|
||||
func TestSync_DeletesDisappearedFiles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
require.NoError(t, os.MkdirAll(filepath.Join(dir, "wiki"), 0o755))
|
||||
// store has a path that doesn't exist on disk anymore
|
||||
store := &stubStore{known: map[string]struct{}{"wiki/old/facts/ghost.md": {}}}
|
||||
res, err := vectorstore.Sync(context.Background(), dir, &stubStoreWithDelete{stubStore: store}, stubEmbedder{vec: make([]float32, 768)})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, res.Deleted)
|
||||
}
|
||||
|
||||
// stubStoreWithDelete is a thin wrapper to capture Delete calls;
|
||||
// stubStore already implements Delete but we need the wrapper to mix
|
||||
// store interfaces with sync-specific expectations.
|
||||
type stubStoreWithDelete struct {
|
||||
*stubStore
|
||||
}
|
||||
|
||||
func TestSync_SkipsIndexFiles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
writeNote(t, dir, "wiki/a/_index.md", "moc")
|
||||
writeNote(t, dir, "wiki/a/facts/real.md", "body")
|
||||
|
||||
store := &stubStore{known: map[string]struct{}{}}
|
||||
res, err := vectorstore.Sync(context.Background(), dir, store, stubEmbedder{vec: make([]float32, 768)})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, res.Added)
|
||||
assert.NotContains(t, store.upserts, "wiki/a/_index.md")
|
||||
}
|
||||
|
||||
func TestSync_NoOpWhenComponentsNil(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
writeNote(t, dir, "wiki/a/facts/x.md", "x")
|
||||
res, err := vectorstore.Sync(context.Background(), dir, nil, nil)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 0, res.Added)
|
||||
}
|
||||
|
||||
func TestSync_CollectsEmbedderErrors(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
writeNote(t, dir, "wiki/a/facts/x.md", "x")
|
||||
store := &stubStore{known: map[string]struct{}{}}
|
||||
emb := stubEmbedder{err: errors.New("upstream down")}
|
||||
res, err := vectorstore.Sync(context.Background(), dir, store, emb)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 0, res.Added)
|
||||
assert.Len(t, res.Errors, 1)
|
||||
}
|
||||
Reference in New Issue
Block a user