package eino import ( "context" "encoding/json" "strings" "github.com/cloudwego/eino/components/retriever" "github.com/cloudwego/eino/schema" "github.com/sundynix/sundynix-shared/contract" ) const defaultRetrieveTopK = 4 // ragRetriever 把 mcp-go 的混合检索(NATS kb_search)包成 Eino components/retriever.Retriever。 // kb 在构造时绑定(owner 作用域库名);查询经 NATS 调 mcp-go,命中转 *schema.Document。 type ragRetriever struct { caller ToolCaller kb string topK int } // newRetriever 构造一个绑定到指定库的 Eino Retriever 组件。 func (o *Orchestrator) newRetriever(kb string) retriever.Retriever { return &ragRetriever{caller: o.tools, kb: kb, topK: defaultRetrieveTopK} } // Retrieve 实现 retriever.Retriever:经 NATS kb_search 检索,返回命中文档(不可用时降级空)。 func (r *ragRetriever) Retrieve(ctx context.Context, query string, _ ...retriever.Option) ([]*schema.Document, error) { if r.caller == nil || r.kb == "" { return nil, nil } topK := r.topK if topK <= 0 { topK = defaultRetrieveTopK } cctx, cancel := context.WithTimeout(ctx, toolCallTimeout) defer cancel() res, err := r.caller.CallTool(cctx, contract.ToolSubjectGo("kb_search"), &contract.ToolCall{ Tool: "kb_search", Args: map[string]any{"kb": r.kb, "q": query, "topK": topK}, }) if err != nil || res == nil || !res.OK || res.Content == "" || res.Content == "[]" { return nil, nil } var hits []struct { Text string `json:"text"` Score float64 `json:"score"` } if json.Unmarshal([]byte(res.Content), &hits) != nil { // 非结构化命中:整体作为单篇文档返回(与旧行为对齐)。 return []*schema.Document{{Content: res.Content}}, nil } docs := make([]*schema.Document, 0, len(hits)) for _, h := range hits { docs = append(docs, &schema.Document{ Content: strings.TrimSpace(h.Text), MetaData: map[string]any{"score": h.Score}, }) } return docs, nil }