Files
sundynix-agentix/sundynix-dispatcher/internal/eino/compose_compiler_test.go
T
Blizzard ef6f525a74 refactor(dispatcher): T1.1 退役 graph.go —— 编排引擎收成单一 compose
compose 已默认数天、HITL/多智能体/评测全在其上真实跑过,soak 充分。删掉自研拓扑
解释器这第二套引擎,消灭"双实现 drift"税:

- 删 runGraph(graph.go 的自研解释器)+ composeEnabled/EINO_COMPOSE 逃生舱开关
  + runConversation(仅 runGraph 用的死代码)。
- executeGraph 直接走 runComposeGraph;compose 编译失败兜底改单轮对话(不再回退
  graph.go);清掉仅 runGraph 用的 import(otel attribute/trace/otelx)。
- 保留 board / 各节点执行器(retriever/tool/agent/branch/approval/map/render/aggregate)
  / 工具函数 —— 它们是 compose 各节点 lambda 复用的,非 graph.go 专属。
- 测试:8 处 runGraph→runComposeGraph;等价测试(对照两引擎)转为 compose 正确性测试;
  runConversation 的开关测试转为「无 ChatModel 降级 runAgent」。

go test ./... 全绿 + vet 干净;冒烟 简单 agent/分支图 跑通。此后每个编排改动不再两边
对齐,成本减半。DEPTH_ROADMAP T1.1 。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-30 08:53:02 +08:00

286 lines
12 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"
"encoding/json"
"errors"
"strings"
"testing"
"time"
"github.com/sundynix/sundynix-dispatcher/internal/harness"
"github.com/sundynix/sundynix-dispatcher/internal/llm"
"github.com/sundynix/sundynix-shared/contract"
)
// echoLLM 回显最后一条 user 消息内容(确定性),便于两条执行路径逐字对比。
func echoLLM() *fakeLLM {
return &fakeLLM{
ready: true,
stream: func(m []llm.ChatMessage) string {
for i := len(m) - 1; i >= 0; i-- {
if m[i].Role == "user" {
return "ANS:" + m[i].Content
}
}
return "ANS:"
},
}
}
// runCompose 跑一个图并返回成稿(自研 graph.go 已退役,编排引擎统一为 compose)。
func runCompose(t *testing.T, graph string) string {
t.Helper()
o := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(), sink: &fakeSink{}}
ans, _, err := o.runComposeGraph(context.Background(), &contract.Task{ID: "t_eq", Graph: []byte(graph)}, &execTracer{})
if err != nil {
t.Fatalf("runComposeGraph: %v", err)
}
return ans
}
// rejectWaiter 是恒拒绝的 HITL 审批替身(fail-safe 路径回归用)。
type rejectWaiter struct{ note string }
func (r *rejectWaiter) WaitApproval(context.Context, string, time.Duration) (*contract.ApprovalDecision, error) {
return &contract.ApprovalDecision{Approved: false, Note: r.note}, nil
}
// TestComposeApprovalRejectStopsDownstream 钉死 compose 路径的 HITL 审批:
// 审批拒绝须置 errRejected、产出拒绝语,且下游 agent 绝不执行(否则等于审批形同虚设)。
// 这是翻默认前最高危的回归点——execDSLNode 一度没有 approval 分支会把审批节点当未识别跳过。
func TestComposeApprovalRejectStopsDownstream(t *testing.T) {
graph := `{"version":"1","nodes":[
{"id":"in","kind":"input","config":{"text":"hi"}},
{"id":"ap","kind":"approval","config":{"title":"上线审批"}},
{"id":"a","kind":"agent","config":{"system":"机密操作"}}
],"edges":[
{"source":"in","target":"ap"},{"source":"ap","target":"a"}
]}`
o := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(),
sink: &fakeSink{}, approval: &rejectWaiter{note: "不批"}}
ans, refs, err := o.runComposeGraph(context.Background(),
&contract.Task{ID: "t_ap", Graph: []byte(graph)}, &execTracer{})
if !errors.Is(err, errRejected) {
t.Fatalf("审批拒绝应返回 errRejectedgot %v", err)
}
if !strings.Contains(ans, "已被拒绝") {
t.Fatalf("应产出拒绝语,got %q", ans)
}
if strings.Contains(ans, "ANS:") {
t.Fatalf("审批拒绝后下游 agent 不应执行,但成稿含 agent 产出: %q", ans)
}
if refs != nil {
t.Fatalf("拒绝终态 refs 应为 nilgot %v", refs)
}
}
// TestComposeApprovalInterruptCheckpoints 钉死中断式 HITL(接了 checkpoint 后端):
// 审批节点应 compose.Interrupt —— 返回 errInterrupted、置 waiting、checkpoint 落盘、下游 agent 不执行。
// 这是「释放 goroutine + 抗重启」的中断半边;恢复半边(resume)由增量3 覆盖。
func TestComposeApprovalInterruptCheckpoints(t *testing.T) {
graph := `{"version":"1","nodes":[
{"id":"in","kind":"input","config":{"text":"hi"}},
{"id":"ap","kind":"approval","config":{"title":"上线审批"}},
{"id":"a","kind":"agent","config":{"system":"机密操作"}}
],"edges":[
{"source":"in","target":"ap"},{"source":"ap","target":"a"}
]}`
kv := newMemKV()
st := &fakeStatus{}
fs := &fakeSink{}
o := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(), sink: fs, status: st}
o.SetCheckpoints(kv) // 接 checkpoint 后端 → 审批走中断模型
task := &contract.Task{ID: "t_intr", Graph: []byte(graph)}
_, _, err := o.runComposeGraph(context.Background(), task, &execTracer{})
if !errors.Is(err, errInterrupted) {
t.Fatalf("审批应中断并返回 errInterruptedgot %v", err)
}
if st.last() != contract.TaskWaiting {
t.Fatalf("中断应置 waitinggot %q", st.last())
}
if _, ok, _ := kv.Get(context.Background(), task.ID); !ok {
t.Fatalf("compose 应把图状态落进 checkpointkey=task_id),但 KV 未命中")
}
if strings.Contains(fs.text(), "ANS:") {
t.Fatalf("审批中断后下游 agent 不应执行,但 sink 含 agent 产出: %q", fs.text())
}
}
// approvalGraph 是 input→approval→agent 的最小 HITL 图,agent 用 echoLLM 回 "ANS:<query>"
// 借此判断审批放行后下游是否真的执行。
const approvalGraph = `{"version":"1","nodes":[
{"id":"in","kind":"input","config":{"text":"写上线报告"}},
{"id":"ap","kind":"approval","config":{"title":"上线审批"}},
{"id":"a","kind":"agent","config":{"system":"报告助手"}}
],"edges":[
{"source":"in","target":"ap"},{"source":"ap","target":"a"}
]}`
// interruptThenSetup 跑一次 fresh,断言中断并落了 resume 记录,返回可继续 resume 的编排器与任务。
func interruptThenSetup(t *testing.T, kv *memKV, st *fakeStatus) (*Orchestrator, *contract.Task) {
t.Helper()
o := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(), sink: &fakeSink{}, status: st}
o.SetCheckpoints(kv)
task := &contract.Task{ID: "t_resume", Graph: []byte(approvalGraph)}
if _, _, err := o.runComposeGraph(context.Background(), task, &execTracer{}); !errors.Is(err, errInterrupted) {
t.Fatalf("fresh 跑应中断,got %v", err)
}
if _, ok, _ := kv.Get(context.Background(), pendingKey(task.ID)); !ok {
t.Fatalf("中断应落 resume 记录(pending:%s)", task.ID)
}
return o, task
}
// TestComposeApprovalResumeApprove 钉死 HITL 恢复半边——批准:从 checkpoint 重入图,黑板无损还原,
// 下游 agent 执行出稿,终态清 resume 记录 + checkpoint。这条用例也透过 eino 真实 checkpoint 路径
// 验证了 board 序列化(增量2a)端到端可用。
func TestComposeApprovalResumeApprove(t *testing.T) {
kv, st := newMemKV(), &fakeStatus{}
o, task := interruptThenSetup(t, kv, st)
dec := &contract.ApprovalDecision{TaskID: task.ID, Approved: true, Note: "放行"}
ans, _, err := o.ResumeApproval(context.Background(), task, dec, &execTracer{})
if err != nil {
t.Fatalf("批准 resume 不应报错: %v", err)
}
if !strings.Contains(ans, "ANS:") {
t.Fatalf("批准后下游 agent 应执行并出稿,got %q", ans)
}
if st.last() != contract.TaskRunning {
t.Fatalf("批准应把状态从 waiting 拨回 runninggot %q", st.last())
}
if _, ok, _ := kv.Get(context.Background(), pendingKey(task.ID)); ok {
t.Fatalf("成功收尾应清 resume 记录")
}
if _, ok, _ := kv.Get(context.Background(), task.ID); ok {
t.Fatalf("成功收尾应清 compose checkpoint")
}
}
// TestComposeApprovalResumeReject 钉死 HITL 恢复半边——拒绝:resume 置 rejected、返回 errRejected、
// 出拒绝语、下游 agent 不执行、清 checkpoint。
func TestComposeApprovalResumeReject(t *testing.T) {
kv, st := newMemKV(), &fakeStatus{}
o, task := interruptThenSetup(t, kv, st)
dec := &contract.ApprovalDecision{TaskID: task.ID, Approved: false, Note: "不批"}
ans, refs, err := o.ResumeApproval(context.Background(), task, dec, &execTracer{})
if !errors.Is(err, errRejected) {
t.Fatalf("拒绝 resume 应返回 errRejectedgot %v", err)
}
if !strings.Contains(ans, "已被拒绝") {
t.Fatalf("应出拒绝语,got %q", ans)
}
if strings.Contains(ans, "ANS:") {
t.Fatalf("拒绝后下游 agent 不应执行,got %q", ans)
}
if refs != nil {
t.Fatalf("拒绝终态 refs 应为 nilgot %v", refs)
}
if _, ok, _ := kv.Get(context.Background(), task.ID); ok {
t.Fatalf("拒绝收尾应清 compose checkpoint")
}
}
// TestHandleApprovalDecisionResumes 钉死 3b 决定驱动恢复:审批决定到达 → 据 resume 记录续跑 →
// 收尾流 + 判 done + 清记录。这是「持久消费者收决定 → ResumeApproval → finishResumed」的闭环。
func TestHandleApprovalDecisionResumes(t *testing.T) {
kv, st, fs := newMemKV(), &fakeStatus{}, &fakeSink{}
o := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(), sink: fs, status: st}
o.SetCheckpoints(kv)
task := &contract.Task{ID: "t_dec", Graph: []byte(approvalGraph)}
if _, _, err := o.runComposeGraph(context.Background(), task, &execTracer{}); !errors.Is(err, errInterrupted) {
t.Fatalf("fresh 跑应中断,got %v", err)
}
// 决定到达(批准)→ 消费者回调驱动续跑 + 收尾。
o.HandleApprovalDecision(context.Background(), &contract.ApprovalDecision{TaskID: task.ID, Approved: true})
if !fs.done {
t.Fatalf("续跑成功应收尾 Token 流")
}
if !strings.Contains(fs.text(), "ANS:") {
t.Fatalf("批准后应续跑下游 agentsink=%q", fs.text())
}
if st.last() != contract.TaskDone {
t.Fatalf("续跑成功应判 donegot %q", st.last())
}
if _, ok, _ := kv.Get(context.Background(), pendingKey(task.ID)); ok {
t.Fatalf("成功收尾应清 resume 记录")
}
}
// TestHandleApprovalDecisionNoPendingIgnored 钉死幂等/兼容:无 resume 记录(阻塞态任务的决定 /
// 决定重投)→ 直接忽略,不收尾、不改状态。
func TestHandleApprovalDecisionNoPendingIgnored(t *testing.T) {
kv, st, fs := newMemKV(), &fakeStatus{}, &fakeSink{}
o := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(), sink: fs, status: st}
o.SetCheckpoints(kv)
o.HandleApprovalDecision(context.Background(), &contract.ApprovalDecision{TaskID: "never_interrupted", Approved: true})
if fs.done || fs.text() != "" {
t.Fatalf("无 resume 记录不应收尾或产出,sink done=%v text=%q", fs.done, fs.text())
}
if seq, _ := st.snapshot(); len(seq) != 0 {
t.Fatalf("无 resume 记录不应改状态,got %v", seq)
}
}
// TestComposeReturnsRefs 钉死 compose 路径回传检索来源——曾被硬写成 nil,导致忠实度评测静默失效。
// 与 runGraph 同图对照:两路径都应回流含检索片段的 refs(喂 grounded judge)。
func TestComposeReturnsRefs(t *testing.T) {
g := `{"nodes":[
{"id":"i","kind":"input","config":{"text":"介绍杭州西湖"}},
{"id":"r","kind":"retriever","config":{"kb":"travel"}},
{"id":"a","kind":"agent","config":{"system":"你是导游"}}
],"edges":[{"source":"i","target":"r"},{"source":"r","target":"a"}]}`
ll := &fakeLLM{ready: true, stream: func(m []llm.ChatMessage) string { return m[0].Content }}
o := newOrch(ll, kbTool(nil), &fakeSink{}, &fakeExec{})
tk := &contract.Task{ID: "t_refs", Graph: json.RawMessage(g), Meta: map[string]any{contract.MetaUserID: "u42"}}
_, refs, err := o.runComposeGraph(context.Background(), tk, &execTracer{})
if err != nil {
t.Fatalf("runComposeGraph: %v", err)
}
if len(refs) == 0 || !strings.Contains(strings.Join(refs, ""), ragSnippet) {
t.Fatalf("compose 路径应回流含检索片段的 refsgot %v", refs)
}
}
// TestComposeLinear 多节点线性图:input→memory→agent→outputcompose 应跑通并产出成稿。
func TestComposeLinear(t *testing.T) {
graph := `{"version":"1","nodes":[
{"id":"in","kind":"input","config":{"text":"什么是图编排"}},
{"id":"m","kind":"memory","config":{}},
{"id":"a","kind":"agent","config":{"system":"你是助手"}},
{"id":"out","kind":"output","config":{}}
],"edges":[
{"source":"in","target":"m"},{"source":"m","target":"a"},{"source":"a","target":"out"}
]}`
if ans := runCompose(t, graph); !strings.Contains(ans, "ANS:") {
t.Fatalf("线性图 compose 应产出成稿,got %q", ans)
}
}
// TestComposeBranch 分支图:input→branch→(真)A/(假)B,条件恒真走 Acompose 应跑通产出。
// 分支真假边的精确选路另由 TestRunGraph_BranchRoutingcompose)钉死。
func TestComposeBranch(t *testing.T) {
graph := `{"version":"1","nodes":[
{"id":"in","kind":"input","config":{"text":"hi"}},
{"id":"br","kind":"branch","config":{"condition":""}},
{"id":"a","kind":"agent","config":{"system":"A"}},
{"id":"b","kind":"agent","config":{"system":"B"}}
],"edges":[
{"source":"in","target":"br"},
{"source":"br","target":"a","sourceHandle":"true"},
{"source":"br","target":"b","sourceHandle":"false"}
]}`
if ans := runCompose(t, graph); !strings.Contains(ans, "ANS:") {
t.Fatalf("分支图 compose 应产出成稿,got %q", ans)
}
}