b26fe21408
排查"向量路为什么是空的"花了半小时,因为空就是空,没有任何线索。这次把
整条检索链上的静默降级一次清掉。
真 bug(不只是可观测性):
- kb_search 与 Search() 都拿 rag.Ready() 当总闸,而 Ready() 只代表"向量路
可用"(embedding + Milvus)。全文(bleve)与图谱(Neo4j)根本不依赖它们,却
被一并毙掉 → "模型配置没下发"表现为"整个知识库什么都搜不到",还不报错。
改为逐路判定,任一路可用就仍有召回。
不再吞错:
- milvus.search 原先把 error 转成 nil,nil —— 检索失败与无召回彻底无法区分;
- bleve.search / graph.search 出错直接回 nil,连日志都没有;
- searchPaths 丢掉 embedding 的 error。
三处改为如实返回,错误统一打日志。
逐路诊断(RouteDiag):每路上报 ok/empty/disabled/error + 耗时 + 原因,经
kb_search 的 diag 参数(仅试验台传,生产调用返回值不变)→ gateway → 检索
试验台。界面上现在能直接看出"这一路没配置/报错了/确实没匹配",不必翻日志。
内存兜底索引也会在 note 里点明"重启即清零"。
测试:3 组,覆盖"无 embedding 时全文仍可召回"、三种空的区分、内存索引提示。
把总闸加回去验证过第一条确实会红——测试能抓到这个回归,不是摆设。
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
441 lines
15 KiB
Go
441 lines
15 KiB
Go
// Package rag 实现 RAG 核心链:embedding(provider 抽象) + Milvus 向量库 + 入库/检索。
|
||
// 是 LLM Wiki 混合检索的向量路;Bleve/Neo4j 融合为后续扩展。
|
||
package rag
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"strconv"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/sundynix/sundynix-shared/contract"
|
||
)
|
||
|
||
const (
|
||
embedBatch = 10 // 每批向量化的块数(兼顾进度可观测 + provider 单批上限)
|
||
embedConcurrency = 4 // 并发批数:几十万字上千块串行太慢,限并发跑快数倍且不打爆 provider 速率
|
||
)
|
||
|
||
// 图谱抽取窗口化参数(env 可配,便于按文档体量/成本调):取代"整篇喂 LLM"(几十万字必爆上下文)。
|
||
// 把已切的语义块合并成 ~graphWindowRunes 的窗口,逐窗并发抽三元组;实体按 kb+name 在 Neo4j 去重。
|
||
func graphWindowRunes() int { return envInt("GRAPH_WINDOW_RUNES", 4000) } // 每窗字数(远小于模型上下文)
|
||
func graphMaxWindows() int { return envInt("GRAPH_MAX_WINDOWS", 60) } // 封顶窗口数(控成本;大文档超出略过+告警)
|
||
func graphConcurrency() int { return envInt("GRAPH_CONCURRENCY", 3) } // 并发抽取窗口数
|
||
|
||
// envInt 读正整数 env,缺省/非法时返回 def。
|
||
func envInt(key string, def int) int {
|
||
if v := os.Getenv(key); v != "" {
|
||
if n, err := strconv.Atoi(v); err == nil && n > 0 {
|
||
return n
|
||
}
|
||
}
|
||
return def
|
||
}
|
||
|
||
// Config 是 RAG 引擎的初始化配置。
|
||
type Config struct {
|
||
MilvusAddr string
|
||
EmbedBase, EmbedKey, EmbedModel string
|
||
RerankBase, RerankKey, RerankModel string
|
||
Neo4jURI, Neo4jUser, Neo4jPass string
|
||
}
|
||
|
||
// Engine 聚合 embedding + Milvus(向量) + Bleve(全文) + Neo4j(图谱) → RRF 融合 + 可选 rerank。
|
||
// embedding 与 chat(图谱抽取用)可热更新(控制面下发)。
|
||
type Engine struct {
|
||
mu sync.RWMutex
|
||
emb *embedClient
|
||
chat *chatClient // 控制面下发的主对话模型
|
||
graphChat *chatClient // 可选:图谱抽取专用便宜模型(env GRAPH_CHAT_*);未配则回退主 chat
|
||
mv *milvusStore
|
||
bleve *bleveStore
|
||
rerank *rerankClient
|
||
graph *graphStore
|
||
}
|
||
|
||
// SetEmbedding 热更新 embedding 配置(控制面变更时调用)。空配置=关闭向量检索。
|
||
func (e *Engine) SetEmbedding(base, key, model string) {
|
||
e.mu.Lock()
|
||
defer e.mu.Unlock()
|
||
if base == "" || model == "" {
|
||
e.emb = nil
|
||
return
|
||
}
|
||
e.emb = newEmbedClient(base, key, model)
|
||
log.Printf("[rag] embedding 配置: %s model=%s", base, model)
|
||
}
|
||
|
||
// SetChat 热更新对话模型配置(图谱实体抽取用,复用控制面 chat 模型)。
|
||
func (e *Engine) SetChat(base, key, model string) {
|
||
e.mu.Lock()
|
||
defer e.mu.Unlock()
|
||
e.chat = newChatClient(base, key, model)
|
||
if e.chat.ready() {
|
||
log.Printf("[rag] 图谱抽取模型: %s model=%s", base, model)
|
||
}
|
||
}
|
||
|
||
func (e *Engine) embed() *embedClient {
|
||
e.mu.RLock()
|
||
defer e.mu.RUnlock()
|
||
return e.emb
|
||
}
|
||
|
||
// Embed 导出当前 embedding 能力供 memory 包复用(满足 memory.Embedder)。
|
||
// 未配置时返回错误,调用方(记忆 Relevance)据此优雅回落。热更新下发的模型即时生效。
|
||
func (e *Engine) Embed(ctx context.Context, texts []string) ([][]float32, error) {
|
||
return e.embed().Embed(ctx, texts)
|
||
}
|
||
|
||
func (e *Engine) chatClient() *chatClient {
|
||
e.mu.RLock()
|
||
defer e.mu.RUnlock()
|
||
return e.chat
|
||
}
|
||
|
||
// graphChatClient 返回图谱抽取用的模型:优先专用便宜模型(env GRAPH_CHAT_*),未配则回退主 chat。
|
||
// 让"每篇最多 60 次图谱抽取"走便宜模型,主对话仍用好模型——成本与质量解耦。
|
||
func (e *Engine) graphChatClient() *chatClient {
|
||
e.mu.RLock()
|
||
defer e.mu.RUnlock()
|
||
if e.graphChat.ready() {
|
||
return e.graphChat
|
||
}
|
||
return e.chat
|
||
}
|
||
|
||
// Open 建立 RAG 引擎。各路连不上 → 降级(不阻断工具服务)。
|
||
func Open(ctx context.Context, cfg Config) *Engine {
|
||
e := &Engine{
|
||
bleve: openBleve(),
|
||
rerank: newRerankClient(cfg.RerankBase, cfg.RerankKey, cfg.RerankModel),
|
||
graph: openGraph(ctx, cfg.Neo4jURI, cfg.Neo4jUser, cfg.Neo4jPass),
|
||
}
|
||
// 图谱抽取专用便宜模型(可选,env GRAPH_CHAT_*):未配则图谱抽取回退控制面主 chat。
|
||
if gb := os.Getenv("GRAPH_CHAT_BASE"); gb != "" {
|
||
e.graphChat = newChatClient(gb, os.Getenv("GRAPH_CHAT_KEY"), os.Getenv("GRAPH_CHAT_MODEL"))
|
||
if e.graphChat.ready() {
|
||
log.Printf("[rag] 图谱抽取专用便宜模型: %s model=%s", gb, os.Getenv("GRAPH_CHAT_MODEL"))
|
||
}
|
||
}
|
||
if e.rerank.ready() {
|
||
log.Printf("[rag] rerank: %s model=%s", cfg.RerankBase, cfg.RerankModel)
|
||
}
|
||
if cfg.EmbedBase != "" && cfg.EmbedModel != "" {
|
||
e.SetEmbedding(cfg.EmbedBase, cfg.EmbedKey, cfg.EmbedModel)
|
||
} else {
|
||
log.Println("[rag] embedding 未配置(待控制面下发),向量检索暂降级")
|
||
}
|
||
if cfg.MilvusAddr != "" {
|
||
mv, err := openMilvus(ctx, cfg.MilvusAddr)
|
||
if err != nil {
|
||
log.Printf("[rag] Milvus 不可用,向量检索降级: %v", err)
|
||
} else {
|
||
e.mv = mv
|
||
log.Printf("[rag] Milvus connected %s", cfg.MilvusAddr)
|
||
}
|
||
}
|
||
return e
|
||
}
|
||
|
||
// Triples 返回某 kb 的图谱三元组(供 UI 可视化)。
|
||
func (e *Engine) Triples(ctx context.Context, kb string, limit int) []Triple {
|
||
return e.graph.triples(ctx, kb, limit)
|
||
}
|
||
|
||
// Ready 报告 RAG 是否可用(embedding + Milvus 均就绪)。
|
||
func (e *Engine) Ready() bool { return e.embed().ready() && e.mv != nil }
|
||
|
||
// Status 报告各依赖子系统的就绪情况(供 health 工具 → 控制台健康灯)。
|
||
func (e *Engine) Status() map[string]bool {
|
||
return map[string]bool{
|
||
"milvus": e.mv != nil,
|
||
"neo4j": e.graph.ready(),
|
||
"embedding": e.embed().ready(),
|
||
// 全文路是否落盘持久。false = 退了内存兜底,重启清零 —— 不上报的话这种降级
|
||
// 只表现为"召回变差",看不出故障(曾因此静默坏了很久)。
|
||
"fulltext_disk": e.bleve != nil && e.bleve.persistent,
|
||
}
|
||
}
|
||
|
||
// Ingest 把一段文本切块 → 分批向量化 → 写 Milvus + Bleve,返回块数。
|
||
// doc 非空表示这是某篇文档/笔记(按 doc 先删旧块再写,支持编辑替换,不重复累积)。
|
||
// onProgress 非空时逐阶段/逐批回调进度(用于实时入库监控)。
|
||
func (e *Engine) Ingest(ctx context.Context, kb, doc, text string, onProgress func(contract.IngestEvent)) (int, error) {
|
||
emit := func(ev contract.IngestEvent) {
|
||
if onProgress != nil {
|
||
onProgress(ev)
|
||
}
|
||
}
|
||
if !e.Ready() {
|
||
return 0, errors.New("rag 未配置(需 embedding + Milvus)")
|
||
}
|
||
chunks := chunk(text)
|
||
if len(chunks) == 0 {
|
||
return 0, nil
|
||
}
|
||
emit(contract.IngestEvent{Stage: "切块", Total: len(chunks), Chunks: previews(chunks), Msg: "拆为 " + itoa(len(chunks)) + " 块"})
|
||
|
||
// 并发分批向量化(保序,逐批回报进度)—— 几十万字上千块时比串行快数倍。
|
||
vecs, err := e.embedAll(ctx, chunks, emit)
|
||
if err != nil {
|
||
emit(contract.IngestEvent{Stage: "失败", Error: "向量化: " + err.Error()})
|
||
return 0, err
|
||
}
|
||
|
||
emit(contract.IngestEvent{Stage: "写Milvus", Msg: "向量库写入中"})
|
||
if len(vecs) > 0 {
|
||
e.mv.deleteDoc(ctx, kb, doc, len(vecs[0])) // 编辑/重入库:先清该 doc 旧块
|
||
}
|
||
if err := e.mv.insert(ctx, kb, doc, chunks, vecs); err != nil {
|
||
emit(contract.IngestEvent{Stage: "失败", Error: "写Milvus: " + err.Error()})
|
||
return 0, err
|
||
}
|
||
emit(contract.IngestEvent{Stage: "写Bleve", Msg: "全文索引写入中"})
|
||
e.bleve.deleteDoc(kb, doc)
|
||
_ = e.bleve.index(kb, doc, chunks) // 同步写全文索引(失败不阻断向量入库)
|
||
|
||
// 图谱路:窗口化并发抽实体/关系 → Neo4j(可降级,不阻断向量入库)。
|
||
// 几十万字按 ~graphWindowRunes 分窗逐窗抽,全覆盖;实体在 Neo4j 按 kb+name 去重。
|
||
if e.graph.ready() && e.graphChatClient().ready() {
|
||
e.extractGraphWindowed(ctx, kb, doc, chunks, emit) // doc=file_id:图谱关系按它标源,供级联删
|
||
}
|
||
|
||
return len(chunks), nil
|
||
}
|
||
|
||
// previews 取每块的前若干字作为预览(供 UI 展示拆分情况)。
|
||
func previews(chunks []string) []string {
|
||
out := make([]string, len(chunks))
|
||
for i, c := range chunks {
|
||
r := []rune(c)
|
||
if len(r) > 50 {
|
||
out[i] = string(r[:50]) + "…"
|
||
} else {
|
||
out[i] = c
|
||
}
|
||
}
|
||
return out
|
||
}
|
||
|
||
func itoa(n int) string {
|
||
if n == 0 {
|
||
return "0"
|
||
}
|
||
var b []byte
|
||
for n > 0 {
|
||
b = append([]byte{byte('0' + n%10)}, b...)
|
||
n /= 10
|
||
}
|
||
return string(b)
|
||
}
|
||
|
||
// Search 混合检索:Milvus(向量) + Bleve(全文) + Neo4j(图谱) → RRF 融合 → 可选 rerank → topK。
|
||
// 注意这里**不再**用 Ready() 当总闸:Ready() 只代表"向量路可用",而全文/图谱两路不依赖
|
||
// embedding 与 Milvus。以前一刀切返回空,导致 embedding 配置缺失时整个知识库像是"什么都搜不到",
|
||
// 且无任何错误信息。现在各路独立判定,任一路可用就仍有召回(见 searchPaths 的 RouteDiag)。
|
||
func (e *Engine) Search(ctx context.Context, kb, query string, topK int) ([]Hit, error) {
|
||
if topK <= 0 {
|
||
topK = 5
|
||
}
|
||
fanout := topK * 3
|
||
|
||
vecHits, ftHits, graphHits, _ := e.searchPaths(ctx, kb, query, fanout)
|
||
// RRF 融合(三路,按文本去重)
|
||
cand := rrf([][]Hit{vecHits, ftHits, graphHits}, fanout)
|
||
log.Printf("[rag] hybrid: 向量=%d 全文=%d 图谱=%d → 融合=%d", len(vecHits), len(ftHits), len(graphHits), len(cand))
|
||
|
||
// 可选 rerank:对融合候选重排取 topK
|
||
if e.rerank.ready() && len(cand) > 1 {
|
||
if rr, rerr := e.rerank.rerank(ctx, query, cand, topK); rerr == nil {
|
||
return rr, nil
|
||
} else {
|
||
log.Printf("[rag] rerank 降级(用 RRF 结果): %v", rerr)
|
||
}
|
||
}
|
||
if len(cand) > topK {
|
||
cand = cand[:topK]
|
||
}
|
||
return cand, nil
|
||
}
|
||
|
||
// DeleteDoc 按 file_id 级联删某文档在三库的痕迹:Milvus 向量块 + Bleve 全文块 + Neo4j 关系。
|
||
// DeleteDoc 按 file_id 级联删三库痕迹。三库**全试一遍**(最大化清理,不因一处失败漏删另两处),
|
||
// 聚合各库错误上返。三库删均幂等(删不存在=no-op),故上层可安全重试。
|
||
func (e *Engine) DeleteDoc(ctx context.Context, kb, fileID string) error {
|
||
if fileID == "" {
|
||
return errors.New("file_id 必填")
|
||
}
|
||
var errs []error
|
||
if err := e.mv.deleteByFile(ctx, kb, fileID); err != nil {
|
||
errs = append(errs, fmt.Errorf("向量: %w", err))
|
||
}
|
||
if err := e.bleve.deleteDoc(kb, fileID); err != nil {
|
||
errs = append(errs, fmt.Errorf("全文: %w", err))
|
||
}
|
||
if err := e.graph.deleteByFile(ctx, kb, fileID); err != nil {
|
||
errs = append(errs, fmt.Errorf("图谱: %w", err))
|
||
}
|
||
if len(errs) > 0 {
|
||
return fmt.Errorf("级联删部分失败(幂等,可重试): %w", errors.Join(errs...))
|
||
}
|
||
log.Printf("[rag] 已删除文档痕迹 kb=%s file_id=%s(向量/全文/图谱)", kb, fileID)
|
||
return nil
|
||
}
|
||
|
||
// RouteDiag 是一路召回的诊断。存在的理由:三路里任何一路挂掉都**不会报错**,
|
||
// 只表现为召回变差——"这一路没配置"、"这一路报错了"、"这一路确实没匹配"
|
||
// 在结果上完全一样(都是空数组),运维无从分辨。检索试验台据此告诉人是哪一环坏了。
|
||
type RouteDiag struct {
|
||
Name string `json:"name"` // vector | fulltext | graph
|
||
Status string `json:"status"` // ok | empty | disabled | error
|
||
Hits int `json:"hits"` //
|
||
MS int64 `json:"ms"` // 该路耗时
|
||
Error string `json:"error,omitempty"` // status=error 时的原文
|
||
Note string `json:"note,omitempty"` // 给人看的解释
|
||
}
|
||
|
||
func diagOf(name string, hits []Hit, err error, disabled bool, note string, started time.Time) RouteDiag {
|
||
d := RouteDiag{Name: name, Hits: len(hits), MS: time.Since(started).Milliseconds(), Note: note}
|
||
switch {
|
||
case disabled:
|
||
d.Status = "disabled"
|
||
case err != nil:
|
||
d.Status, d.Error = "error", err.Error()
|
||
case len(hits) == 0:
|
||
d.Status = "empty"
|
||
default:
|
||
d.Status = "ok"
|
||
}
|
||
return d
|
||
}
|
||
|
||
// searchPaths 跑三路召回,返回各路命中 + 各路诊断(供混合融合、离线评测与检索试验台)。
|
||
// 任何一路失败都不阻断其它路,但失败会被如实记录并打日志——绝不静默当成"没召回"。
|
||
func (e *Engine) searchPaths(ctx context.Context, kb, query string, fanout int) (vec, ft, graph []Hit, diags []RouteDiag) {
|
||
// ── 向量路:embedding 与 Milvus 任一环节失败都要区分出来 ──
|
||
t := time.Now()
|
||
var vErr error
|
||
var vNote string
|
||
vDisabled := !e.embed().ready() || e.mv == nil
|
||
if vDisabled {
|
||
vNote = "embedding 未配置或 Milvus 未连接"
|
||
} else {
|
||
vecs, err := e.embed().Embed(ctx, []string{query})
|
||
switch {
|
||
case err != nil:
|
||
vErr, vNote = fmt.Errorf("embedding: %w", err), "查询向量化失败,向量路本次无贡献"
|
||
case len(vecs) == 0:
|
||
vErr, vNote = errors.New("embedding 返回空向量"), "向量化返回空结果"
|
||
default:
|
||
vec, vErr = e.mv.search(ctx, kb, vecs[0], fanout)
|
||
if vErr != nil {
|
||
vNote = "Milvus 检索失败"
|
||
} else if len(vec) == 0 {
|
||
vNote = "该知识库在 Milvus 中没有向量(未入库或集合被重建过)"
|
||
}
|
||
}
|
||
}
|
||
diags = append(diags, diagOf("vector", vec, vErr, vDisabled, vNote, t))
|
||
|
||
// ── 全文路 ──
|
||
t = time.Now()
|
||
ftDisabled := !e.bleve.ready()
|
||
var ftErr error
|
||
ftNote := ""
|
||
if ftDisabled {
|
||
ftNote = "全文索引未就绪"
|
||
} else {
|
||
ft, ftErr = e.bleve.search(kb, query, fanout)
|
||
if !e.bleve.persistent {
|
||
ftNote = "索引为内存兜底(重启已清零,历史文档需重新入库)"
|
||
}
|
||
}
|
||
diags = append(diags, diagOf("fulltext", ft, ftErr, ftDisabled, ftNote, t))
|
||
|
||
// ── 图谱路 ──
|
||
t = time.Now()
|
||
gDisabled := !e.graph.ready()
|
||
var gErr error
|
||
gNote := ""
|
||
if gDisabled {
|
||
gNote = "Neo4j 未连接或未配置"
|
||
} else {
|
||
graph, gErr = e.graph.search(ctx, kb, query, fanout)
|
||
}
|
||
diags = append(diags, diagOf("graph", graph, gErr, gDisabled, gNote, t))
|
||
|
||
for _, d := range diags {
|
||
if d.Status == "error" {
|
||
log.Printf("[rag] ⚠️ %s 路检索失败 kb=%s: %s", d.Name, kb, d.Error)
|
||
}
|
||
}
|
||
return
|
||
}
|
||
|
||
// SearchByModeDiag 与 SearchByMode 同源,额外返回各路诊断(检索试验台用)。
|
||
// 试验台要回答的是"为什么这一路是空的",光有命中数回答不了。
|
||
func (e *Engine) SearchByModeDiag(ctx context.Context, kb, query string, topK int, mode string) ([]Hit, []RouteDiag) {
|
||
if topK <= 0 {
|
||
topK = 5
|
||
}
|
||
fanout := topK * 3
|
||
vec, ft, graph, diags := e.searchPaths(ctx, kb, query, fanout)
|
||
var hits []Hit
|
||
switch mode {
|
||
case "vector":
|
||
hits = vec
|
||
case "fulltext":
|
||
hits = ft
|
||
case "graph":
|
||
hits = graph
|
||
default:
|
||
hits = rrf([][]Hit{vec, ft, graph}, fanout)
|
||
}
|
||
if len(hits) > topK {
|
||
hits = hits[:topK]
|
||
}
|
||
return hits, diags
|
||
}
|
||
|
||
// SearchByMode 按指定模式返回 topK(评测用,纯检索不 rerank,便于公平对比)。
|
||
// mode: vector|fulltext|graph|hybrid(RRF)。
|
||
func (e *Engine) SearchByMode(ctx context.Context, kb, query string, topK int, mode string) []Hit {
|
||
if topK <= 0 {
|
||
topK = 5
|
||
}
|
||
fanout := topK * 3
|
||
vec, ft, graph, _ := e.searchPaths(ctx, kb, query, fanout)
|
||
var hits []Hit
|
||
switch mode {
|
||
case "vector":
|
||
hits = vec
|
||
case "fulltext":
|
||
hits = ft
|
||
case "graph":
|
||
hits = graph
|
||
default: // hybrid
|
||
hits = rrf([][]Hit{vec, ft, graph}, fanout)
|
||
}
|
||
if len(hits) > topK {
|
||
hits = hits[:topK]
|
||
}
|
||
return hits
|
||
}
|
||
|
||
func (e *Engine) Close() {
|
||
if e.mv != nil {
|
||
e.mv.close()
|
||
}
|
||
if e.bleve != nil {
|
||
e.bleve.close() // 落盘版释放锁并刷盘
|
||
}
|
||
e.graph.close(context.Background())
|
||
}
|
||
|
||
// chunk 的实现已移到 chunk.go(递归 + 句界 + 重叠 + rune 安全的语义切块)。
|