Files
sundynix-agentix/sundynix-dispatcher/internal/eino/compose_graph.go
T
Blizzard 03225e31a9 refactor(dispatcher): compose 翻默认 —— 补平三缺口后退役 graph.go 在望
EINO_COMPOSE 默认改为开(留 =0 逃生舱回退权威 graph.go)。翻默认前补平
compose 路径三个会静默吃掉治理能力的缺口:

- HITL 审批:execDSLNode 无 approval 分支 → 审批节点被当未识别跳过、下游照跑;
  补 case 调 approvalNode。
- 忠实度评测:executeGraph 把 refs 硬写 nil → 评测静默失效;runComposeGraph
  改签名回传 refsOf(b)。
- 终态传播:rejected/fatalErr 不传播也不阻断下游 → 任务误判 done-空;加节点
  lambda 入口守卫 + branch cond 守卫 + 终态上抛 errRejected/fatalErr,并让
  runComposeConversation 模型失败置 fatalErr(对齐 graph.go 的 break 语义)。

新增等价测试:审批拒绝停下游、refs 回流。go test ./... 全绿。
soak 无回归后即可物理退役 graph.go。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 11:13:46 +08:00

128 lines
4.6 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 灰度收尾:compose 已与 graph.go 终态对齐(refs / 审批拒绝 / 预算·模型失败 全部传播),
// 默认翻为 compose 运行时;保留逃生舱 EINO_COMPOSE=0 → 回退权威 graph.go。soak 无回归后即可退役 graph.go。
func composeEnabled() bool { return os.Getenv("EINO_COMPOSE") != "0" }
// 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)
// 成本护栏:计入输入 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 由 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())
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()))))
}