package rag import ( "fmt" "hash/fnv" "log" "os" "path/filepath" "github.com/blevesearch/bleve/v2" "github.com/blevesearch/bleve/v2/analysis/analyzer/keyword" "github.com/blevesearch/bleve/v2/analysis/lang/cjk" "github.com/blevesearch/bleve/v2/mapping" "github.com/blevesearch/bleve/v2/search/query" ) // bleveStore 是全文(BM25)检索路。落盘索引(scorch):随 ingest 写入并持久,进程重启不丢, // 与 Milvus(向量)/Neo4j(图谱) 同为持久存储——保证"三路混合"重启后仍是三路。 // 落盘不可用时退回内存索引(功能在、重启会丢),好过全文路全哑。 type bleveStore struct { idx bleve.Index // persistent=false 表示退了内存兜底:功能在但重启即清零。这是**静默降级**—— // 混合检索只会少一路召回、不报错,所以必须上报出去(health → 控制台),否则没人发现。 persistent bool } // bleveMapping:text 字段用 cjk 分词器(bigram,能切中文,否则默认标准分词器把整段中文当一个 // token → 中文全文检索永远 0 命中);kb/doc 用 keyword(不分词,保 TermQuery 精确过滤)。 func bleveMapping() mapping.IndexMapping { text := bleve.NewTextFieldMapping() text.Analyzer = cjk.AnalyzerName kw := bleve.NewTextFieldMapping() kw.Analyzer = keyword.Name doc := bleve.NewDocumentMapping() doc.AddFieldMappingsAt("text", text) doc.AddFieldMappingsAt("kb", kw) doc.AddFieldMappingsAt("doc", kw) im := bleve.NewIndexMapping() im.DefaultMapping = doc return im } // blevePath 返回全文索引落盘目录(env BLEVE_PATH,默认 ./.data/bleve)。 // 生产部署应指向持久卷;开发期落在工作目录下的 .data。 func blevePath() string { if p := os.Getenv("BLEVE_PATH"); p != "" { return p } return ".data/bleve" } // openBleve 打开(或新建)落盘全文索引;落盘失败退回内存索引兜底。 func openBleve() *bleveStore { path := blevePath() idx, err := bleve.Open(path) if err == bleve.ErrorIndexPathDoesNotExist { if mkerr := os.MkdirAll(filepath.Dir(path), 0o755); mkerr != nil { log.Printf("[rag] bleve 目录创建失败,全文路退内存(重启会丢): %v", mkerr) return memBleve() } idx, err = bleve.New(path, bleveMapping()) } if err != nil { log.Printf("[rag] bleve 落盘打开失败,全文路退内存(重启会丢): %v", err) return memBleve() } log.Printf("[rag] bleve 全文索引落盘就绪: %s", path) return &bleveStore{idx: idx, persistent: true} } // memBleve 内存兜底:落盘不可用时退回内存索引(功能在、重启会丢),好过全文路全哑。 func memBleve() *bleveStore { idx, err := bleve.NewMemOnly(bleveMapping()) if err != nil { log.Printf("[rag] bleve 内存兜底也失败,全文路降级: %v", err) return &bleveStore{} } log.Printf("[rag] bleve 退回内存索引(非持久)") return &bleveStore{idx: idx} } func (b *bleveStore) ready() bool { return b != nil && b.idx != nil } // close 关闭索引(落盘版释放锁并刷盘);优雅停机时调。 func (b *bleveStore) close() { if b.ready() { _ = b.idx.Close() } } // index 把 (kb, doc, texts) 写入全文索引(id 含 kb+doc+文本哈希,幂等)。 func (b *bleveStore) index(kb, doc string, texts []string) error { if !b.ready() { return nil } batch := b.idx.NewBatch() for _, t := range texts { id := fmt.Sprintf("%s:%s:%x", kb, doc, fnvHash(t)) if err := batch.Index(id, map[string]any{"text": t, "kb": kb, "doc": doc}); err != nil { return err } } return b.idx.Batch(batch) } // deleteDoc 删除某 (kb, doc) 的全部全文块(笔记重入库前清旧块)。 func (b *bleveStore) deleteDoc(kb, doc string) error { if !b.ready() || doc == "" { return nil } kq := bleve.NewTermQuery(kb) kq.SetField("kb") dq := bleve.NewTermQuery(doc) dq.SetField("doc") req := bleve.NewSearchRequest(bleve.NewConjunctionQuery(kq, dq)) req.Size = 1000 res, err := b.idx.Search(req) if err != nil { return err } batch := b.idx.NewBatch() for _, h := range res.Hits { batch.Delete(h.ID) } return b.idx.Batch(batch) } // search 全文检索(可按 kb 过滤),返回 BM25 排序的命中。 // 错误如实返回:以前吞掉错误只回 nil,全文路挂了看起来就只是"没召回"。 func (b *bleveStore) search(kb, q string, topK int) ([]Hit, error) { if !b.ready() || q == "" { return nil, nil } mq := bleve.NewMatchQuery(q) mq.SetField("text") var qy query.Query = mq if kb != "" { tq := bleve.NewTermQuery(kb) tq.SetField("kb") qy = bleve.NewConjunctionQuery(mq, tq) } req := bleve.NewSearchRequest(qy) req.Size = topK req.Fields = []string{"text"} res, err := b.idx.Search(req) if err != nil { return nil, err } var hits []Hit for _, h := range res.Hits { text, _ := h.Fields["text"].(string) if text != "" { hits = append(hits, Hit{Text: text, Score: float32(h.Score)}) } } return hits, nil } func fnvHash(s string) uint64 { h := fnv.New64a() _, _ = h.Write([]byte(s)) return h.Sum64() }