Files
Blizzard 53f7e172c3 fix: KB 级联删「事务化」—— 失败不再留不可删孤儿(T4.F)
原为 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>
2026-07-06 16:32:25 +08:00

160 lines
4.6 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
}
// bleveMappingtext 字段用 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()
}