From 45ad8528065b324a3935f56247f642fb7943d54d Mon Sep 17 00:00:00 2001 From: Blizzard Date: Tue, 23 Jun 2026 11:20:52 +0800 Subject: [PATCH] =?UTF-8?q?feat(dispatcher):=20Eino=20=E7=BB=84=E4=BB=B6?= =?UTF-8?q?=E5=8C=96=E8=A1=A5=E5=AE=8C=20=E2=80=94=E2=80=94=20=E6=A3=80?= =?UTF-8?q?=E7=B4=A2=20Retriever=20+=20=E6=8F=90=E7=A4=BA=E8=AF=8D=20ChatT?= =?UTF-8?q?emplate?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 补完终态架构剩余两层(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) --- EINO_ADOPTION.md | 3 + sundynix-dispatcher/internal/eino/compile.go | 22 +++++-- .../internal/eino/eino_components.go | 62 +++++++++++++++++++ .../internal/eino/eino_components_test.go | 56 +++++++++++++++++ sundynix-dispatcher/internal/eino/report.go | 24 ++----- 5 files changed, 142 insertions(+), 25 deletions(-) create mode 100644 sundynix-dispatcher/internal/eino/eino_components.go create mode 100644 sundynix-dispatcher/internal/eino/eino_components_test.go diff --git a/EINO_ADOPTION.md b/EINO_ADOPTION.md index c199335..950fa73 100644 --- a/EINO_ADOPTION.md +++ b/EINO_ADOPTION.md @@ -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) diff --git a/sundynix-dispatcher/internal/eino/compile.go b/sundynix-dispatcher/internal/eino/compile.go index 56bd90c..787d624 100644 --- a/sundynix-dispatcher/internal/eino/compile.go +++ b/sundynix-dispatcher/internal/eino/compile.go @@ -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 把工具入参压成一行短预览。 diff --git a/sundynix-dispatcher/internal/eino/eino_components.go b/sundynix-dispatcher/internal/eino/eino_components.go new file mode 100644 index 0000000..7e8b85c --- /dev/null +++ b/sundynix-dispatcher/internal/eino/eino_components.go @@ -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 +} diff --git a/sundynix-dispatcher/internal/eino/eino_components_test.go b/sundynix-dispatcher/internal/eino/eino_components_test.go new file mode 100644 index 0000000..8a05dcc --- /dev/null +++ b/sundynix-dispatcher/internal/eino/eino_components_test.go @@ -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) + } +} diff --git a/sundynix-dispatcher/internal/eino/report.go b/sundynix-dispatcher/internal/eino/report.go index ce15a98..679159f 100644 --- a/sundynix-dispatcher/internal/eino/report.go +++ b/sundynix-dispatcher/internal/eino/report.go @@ -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() }