Files
sundynix-agentix/sundynix-dispatcher/internal/eino/compose_graph.go
T
Blizzard f238ae1455 feat(dispatcher): 编排式多智能体接力 —— 上游 agent 产出沿图传给下游
让画布上串联的多个 agent 真正协作完成一件事(如 检索→撰写→审查 出报告),
而非各自对原 query 重答:
- board 增 agentOut(各上游 agent 产出按序);RunCtx 增 Upstream,buildMessages
  注入"前序协作 agent 的产出,请在此基础上继续"。
- runAgent / runReactAgent / runComposeConversation 三条路径统一:注入 b.agentOut
  作上下文,产出经 recordAgentOutput 入黑板(append agentOut + 设为当前成稿 answer,
  多 agent 时最后一个 agent 产出即最终成品)。

不需要 eino multiagent 组件——编排式协作由现有图引擎 + 上下文传递实现。
测试 TestAgentCollaborationPassesOutput:下游 agent 确见上游产出、成稿=下游产出;
make test-go 全绿。

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

98 lines
3.5 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"
"io"
"os"
"strings"
"time"
"github.com/cloudwego/eino/compose"
"github.com/cloudwego/eino/schema"
"github.com/sundynix/sundynix-dispatcher/internal/harness"
)
// composeEnabled 报告是否启用 compose.Graph 编排路径(Phase C 灰度开关,默认关 → 走自研 graph.go)。
// 并存策略:EINO_COMPOSE=1 时对话主流程改走 Eino compose 运行时,行为对齐后再逐步退役 graph.go。
func composeEnabled() bool { return os.Getenv("EINO_COMPOSE") == "1" }
// runConversation 是对话/模型节点的统一入口:按灰度开关选 compose.Graph 或自研 runAgent。
// 二者对外行为一致(据黑板拼消息 → 流式回流 token → 累计成稿),便于等价回归。
func (o *Orchestrator) runConversation(ctx context.Context, taskID string, b *board, system string, tr *execTracer, node string) {
if composeEnabled() {
o.runComposeConversation(ctx, taskID, b, system, tr, node)
return
}
o.runAgent(ctx, taskID, b, system, tr, node)
}
// runComposeConversation 用 Eino compose.Graph 跑对话主流程:
// START → ChatModel 节点 → END,编译为 Runnable 后流式执行;可观测经 callbacks 桥到 ExecEvent。
// 模型未就绪 / 编译失败时降级回自研 runAgent,保证不回归。
func (o *Orchestrator) runComposeConversation(ctx context.Context, taskID string, b *board, system string, tr *execTracer, node string) {
cm := o.pool.ChatModel()
if cm == nil {
o.runAgent(ctx, taskID, b, system, tr, node) // 无模型 → 走自研路径的降级桩
return
}
g := compose.NewGraph[[]*schema.Message, *schema.Message]()
if err := g.AddChatModelNode("model", cm); err != nil {
tr.info(node, "system", "compose 降级", "建图失败,退回自研路径:"+err.Error())
o.runAgent(ctx, taskID, b, system, tr, node)
return
}
_ = g.AddEdge(compose.START, "model")
_ = g.AddEdge("model", compose.END)
r, err := g.Compile(ctx)
if err != nil {
tr.info(node, "system", "compose 降级", "编译失败,退回自研路径:"+err.Error())
o.runAgent(ctx, taskID, b, system, tr, node)
return
}
rc := &RunCtx{
UserID: b.uid, SessionID: b.sid,
System: firstNonEmpty(system, defaultAgentSystem),
Query: b.query,
Profile: b.profile,
History: b.history,
ToolOut: append(append([]string{}, b.toolOut...), b.refs...),
Upstream: append([]string{}, b.agentOut...), // 前序协作 agent 产出 → 接力
}
msgs, _ := buildMessages(ctx, rc)
t0 := time.Now()
// ChatModel 的 start/end 由 composeTracercallbacks)落轨迹,这里不再手写 emit(归一)。
sr, err := r.Stream(ctx, msgs, compose.WithCallbacks(composeTracer(tr, node)))
if err != nil {
tr.emit(node, "model", "error", "compose 图执行", err.Error(), time.Since(t0).Milliseconds())
return
}
defer sr.Close()
chunks := 0
var produced strings.Builder // 本节点产出(供下游 agent 接力)
for {
chunk, rerr := sr.Recv()
if rerr == io.EOF {
break
}
if rerr != nil {
tr.emit(node, "model", "error", "compose 图执行", rerr.Error(), time.Since(t0).Milliseconds())
return
}
if chunk.Content == "" {
continue
}
safe, _ := harness.RedactSecrets(chunk.Content) // 输出护栏:逐片脱敏
_ = o.sink.PublishToken(taskID, []byte(safe))
produced.WriteString(safe)
chunks++
}
o.recordAgentOutput(b, produced.String())
tr.info(node, "system", "compose 图", fmt.Sprintf("%d 段输出 / %d 字(Eino compose 运行时)", chunks, len([]rune(produced.String()))))
}