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.Graph(Phase C)或自研 graph.go(默认/权威)。 // 返回 (成稿, 检索来源, error);来源供忠实度评测(compose 路径暂不提供来源 → 返回 nil)。 func (o *Orchestrator) executeGraph(ctx context.Context, t *contract.Task, tr *execTracer) (string, []string, error) { if composeEnabled() { ans, err := o.runComposeGraph(ctx, t, tr) return ans, nil, err } 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()) ans, _, gerr := o.runGraph(ctx, t, tr) // 降级路径丢弃 refs(compose 路径暂不评忠实度) return ans, gerr } 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), "未识别节点,跳过") } }