Files
sundynix-agentix/sundynix-dispatcher/internal/eino/compose_compiler.go
T
Blizzard 2e927eca0a feat(dispatcher): Eino 采纳 Phase C 完成 —— 全图 DSL→compose.Graph 编译器
把整张 DSL 图编译为 Eino compose.Graph 执行(编排归一):
- compose_compiler.go:每节点=Lambda,节点体复用 execDSLNode(全节点类型);
  黑板进 compose 本地状态(WithGenLocalState + ProcessState);branch 走
  AddBranch + 状态感知条件(复用 branchNode);边载荷 flowSignal(注册 no-op
  合并支持 fan-in,真实数据全走黑板)。
- WithNodeTriggerMode(AllPredecessor) DAG 模式:无依赖节点并行调度(效率)。
- Handle→executeGraph 按 EINO_COMPOSE 开关选 compose/graph.go;compose 编译
  失败自动降级回 graph.go(安全网)。默认关,graph.go 仍权威。

等价回归:线性图 + 分支图经解释器与 compose 两路径产出逐字一致(单测);
live 多节点分支图 compose 路径 2800 字答复 eval 1.00、FSM done、0 幽灵;
make test-go 全绿。

过渡期 soak 后翻默认到 compose、退役 graph.go。性能后续:编译图按 DSL-hash 缓存。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 11:02:32 +08:00

218 lines
7.4 KiB
Go
Raw 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 eino
import (
"context"
"fmt"
"sync"
"github.com/cloudwego/eino/compose"
"github.com/sundynix/sundynix-dispatcher/internal/dsl"
"github.com/sundynix/sundynix-shared/contract"
)
// flowSignal 是 compose 编排图的边载荷(占位):真实数据全走 compose 本地状态(*board)
// 边只传"该走了"的信号。注册 no-op 合并以支持 fan-in(多分支汇聚到一个节点)。
type flowSignal struct{}
var registerMergeOnce sync.Once
func registerFlowMerge() {
registerMergeOnce.Do(func() {
compose.RegisterValuesMergeFunc(func([]flowSignal) (flowSignal, error) { return flowSignal{}, nil })
})
}
// executeGraph 按灰度开关选编排实现:compose.GraphPhase C)或自研 graph.go(默认/权威)。
func (o *Orchestrator) executeGraph(ctx context.Context, t *contract.Task, tr *execTracer) (string, error) {
if composeEnabled() {
return o.runComposeGraph(ctx, t, tr)
}
return o.runGraph(ctx, t, tr)
}
// runComposeGraph 把 DSL 图编译为 Eino compose.Graph 并执行(Phase C 编排归一):
// 节点体复用现有 execDSLNode;黑板进 compose 本地状态;branch 走 AddBranch
// DAG 触发模式让无依赖节点并行调度(效率)。编译失败即降级回自研 graph.go(安全网)。
func (o *Orchestrator) runComposeGraph(ctx context.Context, t *contract.Task, tr *execTracer) (string, error) {
registerFlowMerge()
flow, ferr := dsl.Parse(t.Graph)
plan := dsl.Compile(t.Graph)
b := &board{
uid: meta(t, contract.MetaUserID),
sid: meta(t, contract.MetaSessionID),
query: plan.Query,
}
// 无图/空图:退化为 compose 单轮对话。
if ferr != nil || flow == nil || len(flow.Nodes) == 0 {
tr.info("task", "system", "无结构化图", "按单轮对话执行(compose")
b.profile = o.fetchMemory(ctx, b.uid, b.query)
b.history = o.fetchHistory(ctx, b.sid)
o.runComposeConversation(ctx, t.ID, b, plan.System, tr, "agent")
return b.answer, nil
}
// 邻接 + 入度(只认两端都存在的边)。
nodeByID := make(map[string]dsl.Node, len(flow.Nodes))
outE := make(map[string][]dsl.Edge)
indeg := make(map[string]int, len(flow.Nodes))
for _, n := range flow.Nodes {
nodeByID[n.ID] = n
indeg[n.ID] = 0
}
for _, e := range flow.Edges {
if _, ok := nodeByID[e.Source]; !ok {
continue
}
if _, ok := nodeByID[e.Target]; !ok {
continue
}
outE[e.Source] = append(outE[e.Source], e)
indeg[e.Target]++
}
// 图里无 memory 节点 → 沿用默认:注入画像+历史(与 graph.go 对齐,避免回归)。
hasMemory := false
for _, n := range flow.Nodes {
if n.Kind == "memory" {
hasMemory = true
break
}
}
if !hasMemory {
b.profile = o.fetchMemory(ctx, b.uid, b.query)
b.history = o.fetchHistory(ctx, b.sid)
}
// 建 compose 图:黑板进本地状态(GenLocalState 闭包持有本任务的 b)。
g := compose.NewGraph[flowSignal, flowSignal](
compose.WithGenLocalState(func(context.Context) *board { return b }),
)
key := func(id string) string { return "n_" + id } // 节点 key 加前缀,避开 START/END 保留字
// 1) 加节点(branch 为 passthrough,路由交给 AddBranch)。
for _, n := range flow.Nodes {
node := n
if node.Kind == "branch" {
_ = g.AddLambdaNode(key(node.ID), compose.InvokableLambda(
func(context.Context, flowSignal) (flowSignal, error) { return flowSignal{}, nil }))
continue
}
_ = g.AddLambdaNode(key(node.ID), compose.InvokableLambda(
func(c context.Context, _ flowSignal) (flowSignal, error) {
perr := compose.ProcessState(c, func(sc context.Context, bd *board) error {
o.execDSLNode(sc, t, node, bd, plan, tr)
return nil
})
return flowSignal{}, perr
}))
}
// 2) 连边。branch 用 AddBranch(条件读 board 选下游);其余直连;终端节点连 END。
for _, n := range flow.Nodes {
node := n
outs := outE[node.ID]
if node.Kind == "branch" {
endNodes := map[string]bool{compose.END: true}
for _, e := range outs {
if _, ok := nodeByID[e.Target]; ok {
endNodes[key(e.Target)] = true
}
}
brn := node
cond := func(c context.Context, _ flowSignal) (map[string]bool, error) {
chosen := map[string]bool{}
_ = compose.ProcessState(c, func(sc context.Context, bd *board) error {
for _, tgt := range o.branchNode(brn, bd, outE[brn.ID], nodeByID, tr) {
chosen[key(tgt)] = true
}
return nil
})
if len(chosen) == 0 {
chosen[compose.END] = true // 没选中任何下游 → 收口到 END,避免悬挂
}
return chosen, nil
}
_ = g.AddBranch(key(node.ID), compose.NewGraphMultiBranch(cond, endNodes))
continue
}
if len(outs) == 0 {
_ = g.AddEdge(key(node.ID), compose.END)
continue
}
for _, e := range outs {
if _, ok := nodeByID[e.Target]; ok {
_ = g.AddEdge(key(node.ID), key(e.Target))
}
}
}
// 3) 入口节点(入度 0)连 START。
for _, n := range flow.Nodes {
if indeg[n.ID] == 0 {
_ = g.AddEdge(compose.START, key(n.ID))
}
}
// 4) 编译(DAG 模式:无依赖节点并行调度)。编译失败 → 降级回自研 graph.go(安全网)。
r, cerr := g.Compile(ctx, compose.WithNodeTriggerMode(compose.AllPredecessor))
if cerr != nil {
tr.info("task", "system", "compose 编译失败", "退回自研 graph.go"+cerr.Error())
return o.runGraph(ctx, t, tr)
}
if _, ierr := r.Invoke(ctx, flowSignal{}); ierr != nil {
tr.info("task", "system", "compose 执行告警", ierr.Error()) // 副作用已落 board;下方按需补一段答复
}
// 图里无 agent 节点(纯工具/检索图)也要出一段答复。
if b.answer == "" {
o.runComposeConversation(ctx, t.ID, b, plan.System, tr, "agent")
}
return b.answer, nil
}
// execDSLNode 执行一个非 branch 的 DSL 节点(compose 编译器用;节点体与 graph.go 一致,
// 区别仅 agent 直接走 runAgent/runReactAgent——compose 已是编排层,不再二次套 compose)。
func (o *Orchestrator) execDSLNode(ctx context.Context, t *contract.Task, n dsl.Node, b *board, plan dsl.Plan, tr *execTracer) {
switch n.Kind {
case "input":
if txt := cstr(n.Config, "text"); txt != "" {
b.query = txt
}
tr.info("input:"+n.ID, "system", labelOf(n, "输入"), truncate(b.query, 80))
case "memory":
if cbool(n.Config, "profile") {
b.profile = o.fetchMemory(ctx, b.uid, b.query)
}
if cbool(n.Config, "history") {
b.history = o.fetchHistory(ctx, b.sid)
}
tr.info("memory:"+n.ID, "memory", labelOf(n, "记忆"),
fmt.Sprintf("画像 %d 字 · 历史 %d 条", len([]rune(b.profile)), len(b.history)))
case "retriever":
o.retrieverNode(ctx, n, b, tr)
case "tool":
o.execToolNode(ctx, t.ID, n, b, tr)
case "agent":
sys := firstNonEmpty(cstr(n.Config, "system"), plan.System)
if cbool(n.Config, "autonomous") {
o.runReactAgent(ctx, t.ID, b, sys, n, tr, "agent:"+n.ID)
} else {
o.runAgent(ctx, t.ID, b, sys, tr, "agent:"+n.ID)
}
case "aggregate":
merged := aggregate(cstr(n.Config, "strategy"), append(append([]string{}, b.refs...), b.toolOut...))
b.refs, b.toolOut = merged, nil
tr.info("aggregate:"+n.ID, "system", labelOf(n, "汇聚"), "策略:"+firstNonEmpty(cstr(n.Config, "strategy"), "拼接"))
case "render":
o.renderNode(ctx, t.ID, n, b, tr)
case "map":
o.mapNode(ctx, t.ID, n, b, tr)
case "output":
tr.info("output:"+n.ID, "system", labelOf(n, "输出"), "目标:"+firstNonEmpty(cstr(n.Config, "target"), "屏幕"))
default:
tr.info(n.Kind+":"+n.ID, "system", labelOf(n, n.Kind), "未识别节点,跳过")
}
}