53f7e172c3
原为 best-effort:三库删失败只 log、MinIO 删错误全吞、PG 照删 → 删一半失败即在 向量/全文/图谱/MinIO 留下「PG 无记录、连 file_id 都查不到」的不可删孤儿。 改为「类事务」(跨库 2PC 不可行,退而求其次:不留不可恢复孤儿 + 失败可见可重试): - milvus.deleteByFile / bleve.deleteDoc / blob.Delete 改返回 error(原 void 吞错) - rag.DeleteDoc 三库全试一遍(最大化清理)+ 聚合错误(原只回 Neo4j 的错);三库删幂等 - gateway KbDeleteDoc 失败闭合:先删依赖存储(三库→MinIO)、PG 最后删; 任一存储删失败 → 不删 PG、返 502「未删除请重试」(保留 file_id 供幂等重试) - 语义翻转:从「总能从列表删掉但留孤儿」→「有孤儿风险就不删、报错可重试」 live:杀 mcp-go→删→502+文档保留;mcp-go 活→删→200+清空(清场僵尸进程后验证) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
160 lines
4.6 KiB
Go
160 lines
4.6 KiB
Go
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
|
||
}
|
||
|
||
// 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}
|
||
}
|
||
|
||
// 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 排序的命中。
|
||
func (b *bleveStore) search(kb, q string, topK int) []Hit {
|
||
if !b.ready() || q == "" {
|
||
return 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
|
||
}
|
||
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
|
||
}
|
||
|
||
func fnvHash(s string) uint64 {
|
||
h := fnv.New64a()
|
||
_, _ = h.Write([]byte(s))
|
||
return h.Sum64()
|
||
}
|