Files
sundynix-agentix/sundynix-dispatcher/internal/eino/compose_compiler_test.go
T
Blizzard 515cf7f87a feat(hitl): 增量3b —— 持久决定投递 + 生产激活(HITL 中断/恢复上线)
把中断/恢复模型在 dispatcher 接线启用,并让审批决定持久化抗离线。至此 HITL
从「阻塞 goroutine 等 5min、core NATS 非持久、不抗重启」升级为「持久化中断 +
决定持久投递 + 从 checkpoint 恢复续跑」。

- contract: 新增审批决定流 StreamApprovals(SUNDYNIX_APPROVALS)/通配
  SubjectApprovalAll/消费者 ConsumerApprovals + checkpoint 桶 BucketCheckpoints。
- bus: EnsureApprovalStream(JetStream 流持久捕获 sundynix.approval.>,MaxAge 24h;
  gateway 现有 nc.Publish 的决定被本流自动捕获,无需改 gateway)+ ConsumeApprovals
  (持久消费者,队列组多副本安全)。
- orchestrator: HandleApprovalDecision(据 task_id 取 resume 记录续跑;无记录则忽略,
  兼容阻塞态任务的决定 + 决定重投幂等)+ finishResumed(收尾对齐 Handle 尾段:
  中断/拒绝/预算/失败/成功+评测落历史)。
- main: 开 checkpoint 存储 + 审批流 + 起决定消费者 → SetCheckpoints 启用中断模型;
  任一步失败优雅降级回阻塞模型;停机 drain 在途 resume。

决定经 JetStream 持久:dispatcher 在决定发出时离线,重连后仍消费到并续跑;任一
dispatcher 副本都能据共享 KV 的 checkpoint+记录恢复(HA)。
测试:决定驱动恢复(续跑下游+判 done+清记录)、无记录忽略(幂等/兼容)。
全模块 go build + go test 绿。

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

293 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:"
},
}
}
func runBoth(t *testing.T, graph string) (interp, comp string) {
t.Helper()
task := &contract.Task{ID: "t_eq", Graph: []byte(graph)}
o1 := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(), sink: &fakeSink{}}
a1, _, err := o1.runGraph(context.Background(), task, &execTracer{})
if err != nil {
t.Fatalf("runGraph: %v", err)
}
o2 := &Orchestrator{pool: echoLLM(), breaker: harness.NewCircuitBreaker(), sink: &fakeSink{}}
a2, _, err := o2.runComposeGraph(context.Background(), task, &execTracer{})
if err != nil {
t.Fatalf("runComposeGraph: %v", err)
}
return a1, a2
}
// 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)
}
}
// TestComposeEquivalentLinear 多节点线性图:input→memory→agent→output,两路径成稿应逐字一致。
func TestComposeEquivalentLinear(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"}
]}`
interp, comp := runBoth(t, graph)
if interp == "" || interp != comp {
t.Fatalf("线性图不等价: interp=%q compose=%q", interp, comp)
}
}
// TestComposeEquivalentBranch 分支图:input→branch→(真)A/(假)B,条件恒真应都走 A,两路径一致。
func TestComposeEquivalentBranch(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"}
]}`
interp, comp := runBoth(t, graph)
if interp == "" || interp != comp {
t.Fatalf("分支图不等价: interp=%q compose=%q", interp, comp)
}
}