feat(dispatcher): Eino 组件化补完 —— 检索 Retriever + 提示词 ChatTemplate
补完终态架构剩余两层(A 收尾): - 检索:新增 ragRetriever(eino_components.go)实现 components/retriever.Retriever, 包 mcp-go kb_search(NATS)、命中转 *schema.Document;retrieve() 改经组件, retriever 节点与报告分章检索统一走它。 - 提示词:buildMessages 改用 prompt.FromMessages + MessagesPlaceholder (FString 仅解析模板串、值原样注入 → JSON 花括号安全),产出消息序列与旧手拼等价。 - 工具:mcpTool 即 InvokableTool(Phase B),用于模型自主调用;静态 tool 节点 确定性裸调用 by design。 测试:buildMessages 形状(含花括号)+ ragRetriever 文档解析单测;make test-go 全绿。 性能:compose 每任务编译实测 ~13µs(基准),相对 LLM 秒级可忽略 → 编译图缓存 判为 premature optimization 暂不做;并行效率已由 Phase C DAG 调度拿到。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -12,6 +12,9 @@
|
||||
- [x] **Phase C · 编排归一**:✅ 全图 `DSL→compose.Graph` 编译器(全节点 + branch + DAG 并行调度)+ callbacks→ExecEvent 归一,等价回归通过(EINO_COMPOSE 灰度开关,默认关;过渡期后 graph.go 退役)
|
||||
- [~] **Phase D · 状态化执行**:✅ 任务生命周期 FSM(已完成)/ ⬜ HITL 中断恢复 / ⬜ 多智能体(按场景)
|
||||
|
||||
**组件化补完(A)**:检索 → `ragRetriever`(`components/retriever.Retriever`,`eino_components.go`);提示词 → `buildMessages` 改用 `prompt.FromMessages`+`MessagesPlaceholder`;工具 → `mcpTool`(`InvokableTool`,模型自主调用)。至此终态架构 8 层中 模型/工具/检索/提示词/编排/智能体(单)/可观测 均已 Eino 组件化;剩 人机交互(中断恢复,按场景)。
|
||||
**性能注记**:compose 每任务编译实测 ~13µs(基准 BenchmarkComposeCompile),相对 LLM 秒级可忽略 → 编译图缓存判定为 premature optimization,暂不做;真正的并行效率已由 Phase C 的 DAG 调度(`AllPredecessor`)拿到。
|
||||
|
||||
---
|
||||
|
||||
## 0. 现状审计(2026-06-22)
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
|
||||
"github.com/cloudwego/eino/components/prompt"
|
||||
"github.com/cloudwego/eino/schema"
|
||||
)
|
||||
|
||||
@@ -20,8 +21,17 @@ type RunCtx struct {
|
||||
ToolOut []string // 工具/检索节点产出(含参考资料)
|
||||
}
|
||||
|
||||
// chatTemplate 是会话消息模板:系统提示词 + 历史占位 + 用户输入。
|
||||
// 用 Eino components/prompt.ChatTemplate(FString 仅解析模板串、值原样注入,故 JSON 花括号安全)。
|
||||
var chatTemplate = prompt.FromMessages(schema.FString,
|
||||
schema.SystemMessage("{system}"),
|
||||
schema.MessagesPlaceholder("history", true),
|
||||
schema.UserMessage("{query}"),
|
||||
)
|
||||
|
||||
// buildMessages 把上下文组装为发给模型的消息序列(系统提示词 + 画像 + 工具产出 + 历史 + 用户输入)。
|
||||
func buildMessages(_ context.Context, rc *RunCtx) ([]*schema.Message, error) {
|
||||
// 系统串的动态拼装(按需注入画像/参考资料)留在 Go 侧;最终经 ChatTemplate 成型。
|
||||
func buildMessages(ctx context.Context, rc *RunCtx) ([]*schema.Message, error) {
|
||||
var sys strings.Builder
|
||||
sys.WriteString(rc.System)
|
||||
if rc.Profile != "" {
|
||||
@@ -33,11 +43,11 @@ func buildMessages(_ context.Context, rc *RunCtx) ([]*schema.Message, error) {
|
||||
sys.WriteString("\n\n以下是工具/检索得到的参考资料:\n")
|
||||
sys.WriteString(strings.Join(rc.ToolOut, "\n---\n"))
|
||||
}
|
||||
msgs := make([]*schema.Message, 0, len(rc.History)+2)
|
||||
msgs = append(msgs, schema.SystemMessage(sys.String()))
|
||||
msgs = append(msgs, rc.History...)
|
||||
msgs = append(msgs, schema.UserMessage(rc.Query))
|
||||
return msgs, nil
|
||||
return chatTemplate.Format(ctx, map[string]any{
|
||||
"system": sys.String(),
|
||||
"history": rc.History,
|
||||
"query": rc.Query,
|
||||
})
|
||||
}
|
||||
|
||||
// previewArgs 把工具入参压成一行短预览。
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
package eino
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/cloudwego/eino/schema"
|
||||
|
||||
"github.com/sundynix/sundynix-shared/contract"
|
||||
)
|
||||
|
||||
// TestBuildMessagesShape 锁定 ChatTemplate 产出的消息序列与旧手拼逻辑等价:
|
||||
// [system(含画像+参考资料)] + [历史...] + [user(query)],且 JSON 花括号不被误解析。
|
||||
func TestBuildMessagesShape(t *testing.T) {
|
||||
rc := &RunCtx{
|
||||
System: "你是助手",
|
||||
Query: "今天几号",
|
||||
Profile: "称呼=Dexter",
|
||||
History: []*schema.Message{schema.UserMessage("上一句"), schema.AssistantMessage("上一答", nil)},
|
||||
ToolOut: []string{`{"k":"v"}`}, // 含花括号,验证 FString 不误解析
|
||||
}
|
||||
msgs, err := buildMessages(context.Background(), rc)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(msgs) != 4 {
|
||||
t.Fatalf("应为 system+2历史+user=4 条,实际 %d", len(msgs))
|
||||
}
|
||||
if msgs[0].Role != schema.System || !strings.Contains(msgs[0].Content, "你是助手") ||
|
||||
!strings.Contains(msgs[0].Content, "称呼=Dexter") || !strings.Contains(msgs[0].Content, `{"k":"v"}`) {
|
||||
t.Fatalf("system 消息不对: %q", msgs[0].Content)
|
||||
}
|
||||
if msgs[3].Role != schema.User || msgs[3].Content != "今天几号" {
|
||||
t.Fatalf("user 消息不对: %q", msgs[3].Content)
|
||||
}
|
||||
}
|
||||
|
||||
// fakeCaller 是 ToolCaller 测试替身,固定返回 kb_search 命中。
|
||||
type fakeCaller struct{ content string }
|
||||
|
||||
func (f *fakeCaller) CallTool(_ context.Context, _ string, _ *contract.ToolCall) (*contract.ToolResult, error) {
|
||||
return &contract.ToolResult{OK: true, Content: f.content}, nil
|
||||
}
|
||||
|
||||
// TestRagRetrieverComponent 验证 Retriever 组件把 kb_search 命中转成 *schema.Document。
|
||||
func TestRagRetrieverComponent(t *testing.T) {
|
||||
o := &Orchestrator{tools: &fakeCaller{content: `[{"text":"片段A","score":0.9},{"text":"片段B","score":0.8}]`}}
|
||||
docs, err := o.newRetriever("u/kb").Retrieve(context.Background(), "q")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(docs) != 2 || docs[0].Content != "片段A" || docs[1].Content != "片段B" {
|
||||
t.Fatalf("文档解析不对: %+v", docs)
|
||||
}
|
||||
}
|
||||
@@ -214,29 +214,15 @@ func (o *Orchestrator) writeSection(ctx context.Context, topic, kb, heading stri
|
||||
return strings.TrimSpace(txt)
|
||||
}
|
||||
|
||||
// retrieve 经 mcp-go kb_search 工具检索知识库,整理为可读参考资料。kb 为空或无召回则返回空。
|
||||
// retrieve 经 Eino Retriever 组件(包 mcp-go kb_search)检索知识库,整理为可读参考资料。
|
||||
func (o *Orchestrator) retrieve(ctx context.Context, kb, query string) string {
|
||||
if o.tools == nil || kb == "" {
|
||||
docs, err := o.newRetriever(kb).Retrieve(ctx, query)
|
||||
if err != nil || len(docs) == 0 {
|
||||
return ""
|
||||
}
|
||||
cctx, cancel := context.WithTimeout(ctx, toolCallTimeout)
|
||||
defer cancel()
|
||||
res, err := o.tools.CallTool(cctx, contract.ToolSubjectGo("kb_search"), &contract.ToolCall{
|
||||
Tool: "kb_search", Args: map[string]any{"kb": kb, "q": query, "topK": 4},
|
||||
})
|
||||
if err != nil || res == nil || !res.OK || res.Content == "" || res.Content == "[]" {
|
||||
return ""
|
||||
}
|
||||
var hits []struct {
|
||||
Text string `json:"text"`
|
||||
Score float64 `json:"score"`
|
||||
}
|
||||
if json.Unmarshal([]byte(res.Content), &hits) != nil {
|
||||
return res.Content
|
||||
}
|
||||
var b strings.Builder
|
||||
for i, h := range hits {
|
||||
fmt.Fprintf(&b, "%d. %s\n", i+1, strings.TrimSpace(h.Text))
|
||||
for i, d := range docs {
|
||||
fmt.Fprintf(&b, "%d. %s\n", i+1, strings.TrimSpace(d.Content))
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user