test(dispatcher): Handle 入口级集成测试 —— 钉死核心编排链 + harness 治理栈协同

此前测试都直打 runGraph/evaluate,绕过真正的任务入口 Handle()——熔断、输入护栏门控、
状态机流转、token 预算、异步评测/用量这些治理逻辑的协同从未在入口级被覆盖。

新增 handle_test.go 6 例(全 fake 依赖,可断言各出口):
- HappyPath:running→done + token 流出/收尾 + 异步评测落库(ok)
- ToolFeedsAgent:图执行→工具→agent 全程经 Handle,工具被调、产出注入
- CircuitBreakerOpen:熔断开 → 快速 failed、不执行图
- BudgetExceeded:meta token_budget=3 → failed + 用量回写带 Exceeded
- GuardrailBlocks:灰区 + LLM 分类器命中 → rejected、不进 running
- EmptyTaskDropped:空 id 直接丢弃、无状态回写

配套 fakeStatus/fakeUsageSink;fakeEvalSink 加锁(Handle 异步评测在 goroutine 写)。
go test -race ./internal/eino 干净,四模块全绿。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Blizzard
2026-06-26 14:10:53 +08:00
parent f926f6fd41
commit 32f50cd983
3 changed files with 251 additions and 6 deletions
+3 -3
View File
@@ -135,7 +135,7 @@ Harness = 围绕 LLM 的可靠性 / 安全 / 质量治理层。4 个组件均为
- **compose 引擎仍灰度**(默认走自研 graph.go);compose 路径下 HITL 审批未实现。
- **OCR 多模态是骨架**`parse_document` 真支持 txt/md/csv/docx/xlsx/文本 PDF,但 MinerU/PaddleOCR 扫描件路径是 TODO 空实现(`mineru.py` 极小)。
- **推理模型已部分适配**:reasoning 模型的 `reasoning_content` 已捕获并 surface 到 exec 轨迹「推理过程」事件(不污染答案);尚未在最终答案区把思维链单独流式呈现给前端。
- **测试偏单元**多为小函数单测;核心链路(图执行、RAG 管线集成/e2e 稀疏;无前端 E2E。
- **测试纵深改善中**小函数单测,已补 `Handle()` 入口级集成测试(熔断/护栏/预算/状态机/工具→agent/异步评测 6 例,-race 干净)+ bus 真实 NATS e2e;剩 RAG 管线集成前端 E2E。
---
@@ -190,7 +190,7 @@ Harness = 围绕 LLM 的可靠性 / 安全 / 质量治理层。4 个组件均为
| 可观测性 | ⭐⭐⭐⭐ | OTel+Prometheus+slog;轨迹流会丢事件 |
| 性能 / 并发 | ⭐⭐⭐½ | 双扩机制具备;规模未实测 |
| 代码质量 | ⭐⭐⭐⭐ | 注释充足、降级完善;KbView 大文件需拆 |
| 测试覆盖 | ⭐⭐⭐ | 单测面广但多为小函数;集成/e2e 稀疏 |
| 测试覆盖 | ⭐⭐⭐½ | 单测面广 + Handle 入口集成 + bus 真实 NATS e2e;剩 RAG 管线集成、前端 E2E |
| CI/CD | ⭐⭐⭐⭐ | CI+Release+更新检查;无签名公证 |
| 运维成熟度 | ⭐⭐ | 单点、无备份灾备、无配额、规模未验证 |
| 生态 / 商业化 | ⭐⭐ | 个人项目;计费未闭环 |
@@ -208,7 +208,7 @@ Harness = 围绕 LLM 的可靠性 / 安全 / 质量治理层。4 个组件均为
| ~~P1~~ 🟡 | 高可用(代码层就绪):调度/工具本就队列组可多副本(实测 2 副本 8 任务 4/4 分摊);**网关改为队列组订阅**eval/usage/status/config),多副本不再重复落库/重复计费(单测+实测去重)。剩 NATS 集群 + 网关 LB + PG/Redis HA 属部署期 | 解单点 |
| P1 | 备份/灾备演练(PG/Milvus/Neo4j | 数据安全 |
| ~~P1~~ ✅ | ~~优雅停机 drain~~:三 Go 服务全覆盖(gateway HTTP Shutdown / dispatcher 在途任务跑完 / mcp-go 在途工具回完),SIGTERM 后等在途至 SHUTDOWN_DRAIN_TIMEOUT(默认30s)再退;在途任务 ctx 脱离信号 ctx 不被掐断。**+ exec 轨迹 Redis 回放**(与 token 流同构,连晚/刷新重连不丢轨迹,实测事后连仍补齐全程)| 可靠性 |
| P1 | 核心链路集成测试 + 前端 E2E | 测试纵深 |
| ~~P1~~ 🟡 | ~~核心链路集成测试~~:Handle 入口级 6 例(熔断/护栏/预算/状态机/工具→agent/异步评测)已补。剩 RAG 管线集成 + 前端 E2E | 测试纵深 |
| P2 | 计量计费闭环 + 按用户配额;OCR 完善;KbView 拆分;循环节点 | |
| P2 | 安全审计 / 渗透(让"纵深防御"从自述变已验证) | 外部 |
| P3 | 多租户/团队;插件体系;K8s Helm;代码签名 | |
@@ -0,0 +1,229 @@
package eino
import (
"context"
"encoding/json"
"sync"
"testing"
"time"
"github.com/sundynix/sundynix-dispatcher/internal/harness"
"github.com/sundynix/sundynix-dispatcher/internal/llm"
"github.com/sundynix/sundynix-shared/contract"
)
// ---- Handle 集成测试:覆盖真正的任务入口(熔断→输入护栏→状态机→执行→评测/用量),
// 而非直接打 runGraph。把核心编排链 + harness 治理栈的协同钉死。----
// fakeStatus 记录任务状态流转序列。
type fakeStatus struct {
mu sync.Mutex
seq []string
detail string
}
func (s *fakeStatus) PublishTaskStatus(_, status, detail string) error {
s.mu.Lock()
s.seq = append(s.seq, status)
s.detail = detail
s.mu.Unlock()
return nil
}
func (s *fakeStatus) snapshot() ([]string, string) {
s.mu.Lock()
defer s.mu.Unlock()
return append([]string{}, s.seq...), s.detail
}
func (s *fakeStatus) last() string {
seq, _ := s.snapshot()
if len(seq) == 0 {
return ""
}
return seq[len(seq)-1]
}
// fakeUsageSink 捕获 token 用量回写(并发安全:Handle 在 defer/goroutine 收尾时发)。
type fakeUsageSink struct {
mu sync.Mutex
ev *contract.UsageEvent
}
func (u *fakeUsageSink) PublishUsage(ev *contract.UsageEvent) error {
u.mu.Lock()
u.ev = ev
u.mu.Unlock()
return nil
}
func (u *fakeUsageSink) get() *contract.UsageEvent {
u.mu.Lock()
defer u.mu.Unlock()
return u.ev
}
// handleKit 是 Handle 测试的依赖套件(全 fake,可断言各出口)。
type handleKit struct {
o *Orchestrator
sink *fakeSink
status *fakeStatus
exec *fakeExec
eval *fakeEvalSink
usage *fakeUsageSink
tools *fakeTools
}
func newHandleKit(ll *fakeLLM) *handleKit {
k := &handleKit{
sink: &fakeSink{}, status: &fakeStatus{}, exec: &fakeExec{},
eval: &fakeEvalSink{}, usage: &fakeUsageSink{}, tools: &fakeTools{},
}
o := &Orchestrator{
pool: ll, breaker: harness.NewCircuitBreaker(),
sink: k.sink, tools: k.tools, exec: k.exec, status: k.status, evalSink: k.eval,
}
o.SetUsageSink(k.usage)
k.o = o
return k
}
// eventually 在 timeout 内轮询 cond,成立即返回 true(等异步评测/用量回写)。
func eventually(timeout time.Duration, cond func() bool) bool {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
if cond() {
return true
}
time.Sleep(10 * time.Millisecond)
}
return cond()
}
const convGraph = `{"nodes":[
{"id":"i","kind":"input","config":{"text":"介绍杭州"}},
{"id":"a","kind":"agent","config":{"system":"助手"}}
],"edges":[{"source":"i","target":"a"}]}`
// happy path:状态 running→donetoken 流出 + 收尾,异步评测落库。
func TestHandle_HappyPath(t *testing.T) {
ll := &fakeLLM{ready: true, stream: func([]llm.ChatMessage) string { return "杭州是浙江省会,以西湖闻名。" }}
k := newHandleKit(ll)
k.o.eval = harness.NewEvaluator(func() bool { return true },
func(context.Context, string, string) (string, error) { return `{"score":5,"reason":"好"}`, nil })
if err := k.o.Handle(context.Background(), task(convGraph)); err != nil {
t.Fatalf("Handle err: %v", err)
}
seq, _ := k.status.snapshot()
if len(seq) < 2 || seq[0] != contract.TaskRunning || seq[len(seq)-1] != contract.TaskDone {
t.Fatalf("状态应 running→…→donegot %v", seq)
}
if !k.sink.done {
t.Error("应收尾 token 流")
}
if k.sink.text() == "" {
t.Error("应有 token 流出")
}
if !eventually(2*time.Second, func() bool { return k.eval.get() != nil }) {
t.Fatal("异步评测应落库")
}
if ev := k.eval.get(); ev.Level != contract.EvalOK {
t.Errorf("评测应 okgot %s", ev.Level)
}
}
// 工具→agent 链路经 Handle 全程跑通:工具被调,产出注入答案。
func TestHandle_ToolFeedsAgent(t *testing.T) {
g := `{"nodes":[
{"id":"i","kind":"input","config":{"text":"hi"}},
{"id":"t","kind":"tool","config":{"tool":"wiki_search"}},
{"id":"a","kind":"agent","config":{"system":"S"}}
],"edges":[{"source":"i","target":"t"},{"source":"t","target":"a"}]}`
ll := &fakeLLM{ready: true, stream: func(m []llm.ChatMessage) string { return m[0].Content }}
k := newHandleKit(ll)
k.tools.fn = func(*contract.ToolCall) *contract.ToolResult {
return &contract.ToolResult{OK: true, Content: "TOOLDATA"}
}
if err := k.o.Handle(context.Background(), task(g)); err != nil {
t.Fatalf("Handle err: %v", err)
}
if !k.tools.called("wiki_search") {
t.Error("应调用 wiki_search")
}
if k.status.last() != contract.TaskDone {
t.Errorf("应 donegot %s", k.status.last())
}
}
// 熔断开启:快速拒绝,不执行,状态 failed。
func TestHandle_CircuitBreakerOpen(t *testing.T) {
ll := &fakeLLM{ready: true, stream: func([]llm.ChatMessage) string { return "不该被调用" }}
k := newHandleKit(ll)
for i := 0; i < 10; i++ { // 连续失败打开熔断
k.o.breaker.Report(false)
}
if k.o.breaker.Allow() {
t.Skip("熔断阈值未达(实现差异),跳过")
}
if err := k.o.Handle(context.Background(), task(convGraph)); err != nil {
t.Fatalf("Handle err: %v", err)
}
if k.status.last() != contract.TaskFailed {
t.Errorf("熔断应 failedgot %s", k.status.last())
}
if k.sink.text() != "" && k.sink.text() != "⚠️ 服务繁忙(已触发熔断保护),请稍后重试。" {
t.Errorf("熔断不应执行图,got token %q", k.sink.text())
}
}
// token 预算触顶:状态 failed(带原因),用量回写带 Exceeded。
func TestHandle_BudgetExceeded(t *testing.T) {
ll := &fakeLLM{ready: true, stream: func([]llm.ChatMessage) string { return "答案" }}
k := newHandleKit(ll)
tk := task(convGraph)
tk.Meta[contract.MetaTokenBudget] = float64(3) // 极小预算,入口计输入即触顶
if err := k.o.Handle(context.Background(), tk); err != nil {
t.Fatalf("Handle err: %v", err)
}
if k.status.last() != contract.TaskFailed {
t.Fatalf("预算超限应 failedgot %s", k.status.last())
}
if !eventually(2*time.Second, func() bool { u := k.usage.get(); return u != nil && u.Exceeded }) {
t.Fatal("应回写用量且标记 Exceeded")
}
}
// 输入护栏 Tier2:灰区任务被 LLM 分类器拦截 → rejected,不执行。
func TestHandle_GuardrailBlocks(t *testing.T) {
ll := &fakeLLM{ready: true, stream: func([]llm.ChatMessage) string { return "不该被调用" }}
k := newHandleKit(ll)
k.o.SetGuardian(harness.NewClassifier(func() bool { return true },
func(context.Context, string, string) (string, error) {
return `{"jailbreak":true,"severity":0.95,"reason":"越狱尝试"}`, nil
}))
tk := task(convGraph)
tk.Meta[contract.MetaSafetyCheck] = true
if err := k.o.Handle(context.Background(), tk); err != nil {
t.Fatalf("Handle err: %v", err)
}
if k.status.last() != contract.TaskRejected {
t.Fatalf("应 rejectedgot %s", k.status.last())
}
seq, _ := k.status.snapshot()
for _, s := range seq {
if s == contract.TaskRunning {
t.Error("被护栏拦截不应进入 running")
}
}
}
// 空 id 任务:直接丢弃,无状态回写(防脏数据自我放大)。
func TestHandle_EmptyTaskDropped(t *testing.T) {
k := newHandleKit(&fakeLLM{ready: true})
if err := k.o.Handle(context.Background(), &contract.Task{ID: "", Graph: json.RawMessage(convGraph)}); err != nil {
t.Fatalf("Handle err: %v", err)
}
if seq, _ := k.status.snapshot(); len(seq) != 0 {
t.Errorf("空任务不应有状态回写,got %v", seq)
}
}
@@ -3,6 +3,7 @@ package eino
import (
"context"
"strings"
"sync"
"testing"
"github.com/sundynix/sundynix-dispatcher/internal/harness"
@@ -10,10 +11,25 @@ import (
"github.com/sundynix/sundynix-shared/contract"
)
// fakeEvalSink 捕获回流的评测事件,供断言纠偏终值。
type fakeEvalSink struct{ ev *contract.EvalEvent }
// fakeEvalSink 捕获回流的评测事件,供断言纠偏终值。并发安全(Handle 异步评测会在 goroutine 写)。
type fakeEvalSink struct {
mu sync.Mutex
ev *contract.EvalEvent
}
func (f *fakeEvalSink) PublishEval(ev *contract.EvalEvent) error { f.ev = ev; return nil }
func (f *fakeEvalSink) PublishEval(ev *contract.EvalEvent) error {
f.mu.Lock()
f.ev = ev
f.mu.Unlock()
return nil
}
// get 并发安全地取最近一次评测事件。
func (f *fakeEvalSink) get() *contract.EvalEvent {
f.mu.Lock()
defer f.mu.Unlock()
return f.ev
}
// orchForEval 组一个仅评测/纠偏所需依赖的编排器(生成器 gen + 评审 judge + 评测出口 sink)。
func orchForEval(gen *fakeLLM, judge func(ctx context.Context, sys, user string) (string, error), sink *fakeEvalSink) *Orchestrator {