c2811e79a3
现象:编排里并排三个 agent(研究/撰写/审查)跑完,运行页没有团队 tab。 两处卡住,不是一处: 1) isMultiAgent 只认 coordinator: 节点。用户自己在图里并排多个 agent 也是 团队,却被整个漏掉。改成:有协调者,或 ≥2 个 agent: 节点。 2) 就算放宽 1,deriveTeam 仍按 kind==="agent" 挑工位——而两条产生 agent 的 路径 kind 并不一致:协调者派发的专家是 kind=agent,图里的 agent 节点是 kind=model。改成按节点名前缀(agent:/tool:)判,这本来就是后端一直遵守的 约定;顺带天然把 retriever:/map:/render: 这些同为 kind=tool 的节点挡在 工位之外(之前它们会混进来当工位)。 连带修一个更要命的:runAgent 把轨迹标签写死成"模型流式推理"、runReactAgent 写死成"ReAct 智能体(自主调工具)",用户在编排里给节点起的名字(研究 Agent / 撰写 Agent / 审查 Agent)整个丢了。后果不止办公室:执行轨迹里三行同名,根本 分不出谁是谁;团队视图只能退回节点 ID,工位显示成 r/w/rev。 改成一律 labelOf(n, 兜底) 由调用方传入,+2 单测钉住。 另:没有协调者时不再凭空画一个"协调者"小人(白板改挂「任务产出」),也不演 递简报那一程(没人可递),✓ 气泡改为收工即冒。 注意:已存的历史轨迹是落库的,仍是旧标签;只有新跑的任务才有节点名。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
112 lines
3.7 KiB
Go
112 lines
3.7 KiB
Go
package eino
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"io"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/cloudwego/eino/compose"
|
||
"github.com/cloudwego/eino/schema"
|
||
|
||
"github.com/sundynix/sundynix-dispatcher/internal/harness"
|
||
)
|
||
|
||
// 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, label string) {
|
||
cm := o.pool.ChatModel()
|
||
if cm == nil {
|
||
o.runAgent(ctx, taskID, b, system, tr, node, label) // 无模型 → 降级桩
|
||
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, label)
|
||
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, label)
|
||
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)
|
||
// 成本护栏:计入输入 token;触顶则中止整图。
|
||
if bud := harness.BudgetFrom(ctx); bud != nil {
|
||
for _, m := range msgs {
|
||
bud.AddPrompt(m.Content)
|
||
}
|
||
if bud.Exceeded() {
|
||
if b.fatalErr == nil {
|
||
b.fatalErr = errBudget
|
||
}
|
||
tr.emit(node, "system", "error", "token 预算", "已达单任务预算上限,中止", 0)
|
||
return
|
||
}
|
||
}
|
||
|
||
t0 := time.Now()
|
||
// ChatModel 的 start/end 由 composeTracer(callbacks)落轨迹,这里不再手写 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())
|
||
if b.fatalErr == nil { // 模型失败 → 标记致命错,让任务判 failed(对齐 runAgent,杜绝 done-空)
|
||
b.fatalErr = fmt.Errorf("agent 模型推理失败: %w", err)
|
||
}
|
||
return
|
||
}
|
||
defer sr.Close()
|
||
|
||
chunks := 0
|
||
var produced strings.Builder // 本节点产出(供下游 agent 接力)
|
||
red := harness.NewStreamRedactor() // 输出护栏:跨分片脱敏,杜绝密钥被切断而漏检
|
||
emit := func(safe string) {
|
||
if safe == "" {
|
||
return
|
||
}
|
||
_ = o.sink.PublishToken(taskID, []byte(safe))
|
||
produced.WriteString(safe)
|
||
chunks++
|
||
}
|
||
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())
|
||
if b.fatalErr == nil { // 流式中断也算失败(结果不完整)→ 判 failed,不静默 done-空
|
||
b.fatalErr = fmt.Errorf("agent 模型推理失败: %w", rerr)
|
||
}
|
||
return
|
||
}
|
||
if chunk.Content == "" {
|
||
continue
|
||
}
|
||
emit(red.Push(chunk.Content))
|
||
}
|
||
emit(red.Flush()) // 吐出暂留尾部
|
||
if bud := harness.BudgetFrom(ctx); bud != nil {
|
||
bud.AddComplete(produced.String()) // 成本护栏:计入输出 token
|
||
}
|
||
o.recordAgentOutput(b, produced.String())
|
||
tr.info(node, "system", "compose 图", fmt.Sprintf("%d 段输出 / %d 字(Eino compose 运行时)", chunks, len([]rune(produced.String()))))
|
||
}
|