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>
This commit is contained in:
Blizzard
2026-06-29 13:08:41 +08:00
parent 368938ce89
commit 515cf7f87a
6 changed files with 217 additions and 2 deletions
@@ -14,6 +14,7 @@ import (
"github.com/sundynix/sundynix-dispatcher/internal/harness"
"github.com/sundynix/sundynix-dispatcher/internal/llm"
dnats "github.com/sundynix/sundynix-dispatcher/internal/nats"
"github.com/sundynix/sundynix-shared/contract"
"github.com/sundynix/sundynix-shared/otelx"
"github.com/sundynix/sundynix-shared/secrets"
)
@@ -57,6 +58,23 @@ func main() {
orch.SetGuardian(guardian) // 输入护栏 Tier2
orch.SetUsageSink(sub) // 成本护栏:token 用量回写网关累计/计费
// HITL 持久化中断/恢复:开 checkpoint 存储 + 审批决定流,审批节点改走中断模型
// compose.Interrupt 落盘释放 goroutine、抗 dispatcher 重启)。任一步失败则降级回阻塞模型。
var drainApprovals func(context.Context)
if kv, kerr := sub.Checkpoints(context.Background(), contract.BucketCheckpoints, 24*time.Hour); kerr != nil {
log.Printf("[dispatcher] 开 checkpoint 存储失败,审批降级为阻塞模型: %v", kerr)
} else if serr := sub.EnsureApprovalStream(context.Background()); serr != nil {
log.Printf("[dispatcher] 建审批决定流失败,审批降级为阻塞模型: %v", serr)
} else {
orch.SetCheckpoints(kv) // 启用中断模型
if drain, cerr := sub.ConsumeApprovals(context.Background(), orch.HandleApprovalDecision); cerr != nil {
log.Printf("[dispatcher] 启动审批决定消费者失败: %v", cerr)
} else {
drainApprovals = drain
log.Println("[dispatcher] HITL 中断/恢复已启用(checkpoint + 审批决定流)")
}
}
// 健康心跳:dispatcher 无 HTTP/工具端点,挂一个 NATS 应答让管理端「服务状态」探到它在线。
startedAt := time.Now()
if unsub, herr := sub.ServeHealth(func() []byte {
@@ -82,6 +100,12 @@ func main() {
if err := sub.ConsumeTasks(ctx, orch.Handle); err != nil && err != context.Canceled {
log.Fatalf("[dispatcher] exit: %v", err)
}
// 优雅停机:停审批决定消费者并等在途 resume 跑完(与任务 drain 同语义)。
if drainApprovals != nil {
dctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
drainApprovals(dctx)
cancel()
}
}
func envOr(key, def string) string {
@@ -5,6 +5,7 @@ import (
"encoding/json"
"errors"
"fmt"
"log"
"sync"
"github.com/cloudwego/eino/compose"
@@ -430,3 +431,69 @@ func (o *Orchestrator) ResumeApproval(ctx context.Context, t *contract.Task, dec
}
return o.execComposeGraph(ctx, t, tr, &resumeCtx{interruptID: p.InterruptID, dec: dec})
}
// HandleApprovalDecision 是审批决定持久消费者的回调:据 task_id 取 resume 记录续跑并收尾。
// 无 resume 记录(阻塞态任务的决定 / 已恢复 / 已过期)则忽略——本消费者只管中断/恢复模型。
// 幂等:决定重投时 resume 记录已被成功收尾清掉 → loadResume 落空 → 安全跳过。
func (o *Orchestrator) HandleApprovalDecision(ctx context.Context, dec *contract.ApprovalDecision) {
if o.checkpoints == nil || dec == nil || dec.TaskID == "" {
return
}
p, err := o.loadResume(ctx, dec.TaskID)
if errors.Is(err, errNoPending) {
return // 非中断态任务(阻塞模型自己经 WaitApproval 收)/ 已处理过
}
if err != nil {
log.Printf("[eino] 取 resume 记录失败 task=%s: %v", dec.TaskID, err)
return
}
var t contract.Task
if uerr := json.Unmarshal(p.Task, &t); uerr != nil || t.ID == "" {
log.Printf("[eino] resume 记录损坏 task=%s: %v", dec.TaskID, uerr)
o.clearResume(ctx, dec.TaskID) // 坏记录清掉,避免决定无限重投
return
}
tr := o.tracer(t.ID)
defer tr.done()
answer, refs, rerr := o.ResumeApproval(ctx, &t, dec, tr)
o.finishResumed(ctx, &t, answer, refs, rerr)
}
// finishResumed 给 resume 续跑收尾,语义对齐 Handle 尾段(中断/拒绝/预算/失败/成功)。
func (o *Orchestrator) finishResumed(ctx context.Context, t *contract.Task, answer string, refs []string, err error) {
switch {
case errors.Is(err, errInterrupted):
return // 续跑中又遇审批:保持 waiting,等下一个决定(resume 记录已重新落盘)
case errors.Is(err, errRejected):
if answer != "" {
_ = o.sink.PublishToken(t.ID, []byte(answer))
}
_ = o.sink.CompleteStream(t.ID)
o.breaker.Report(true) // 拒绝是人为决策,非后端故障
o.setStatus(t.ID, contract.TaskRejected, truncate(answer, 120))
return
case errors.Is(err, errBudget):
if answer != "" {
_ = o.sink.PublishToken(t.ID, []byte(answer))
}
_ = o.sink.PublishToken(t.ID, []byte("\n\n⚠️ 已达单任务 token 预算上限,自动中止。"))
_ = o.sink.CompleteStream(t.ID)
o.breaker.Report(true)
o.setStatus(t.ID, contract.TaskFailed, "token 预算超限")
return
case err != nil:
_ = o.sink.CompleteStream(t.ID)
o.breaker.Report(false)
o.finishStatus(t.ID, err)
return
}
// 成功:收尾流 + 判 done + 评测(低分纠偏)+ 落历史(同 Handle)。
_ = o.sink.CompleteStream(t.ID)
o.breaker.Report(true)
o.finishStatus(t.ID, nil)
go func() {
query := dsl.Compile(t.Graph).Query
final := o.evaluate(t, query, answer, refs)
o.memorize(t, final)
}()
}
@@ -191,6 +191,51 @@ func TestComposeApprovalResumeReject(t *testing.T) {
}
}
// 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) {
@@ -102,6 +102,22 @@ func (s *Subscriber) WaitApproval(ctx context.Context, taskID string, timeout ti
return s.inner.WaitApproval(ctx, taskID, timeout)
}
// Checkpoints 打开 HITL 持久化中断的 KV 桶(compose checkpoint + resume 记录)。
// 返回的 *KVHandle 结构化满足 eino.CheckpointKVGet/Put/Delete),可直接 SetCheckpoints。
func (s *Subscriber) Checkpoints(ctx context.Context, bucket string, ttl time.Duration) (*sharedbus.KVHandle, error) {
return s.inner.Checkpoints(ctx, bucket, ttl)
}
// EnsureApprovalStream 确保审批决定流存在(中断/恢复模型用,决定持久抗 dispatcher 离线)。
func (s *Subscriber) EnsureApprovalStream(ctx context.Context) error {
return s.inner.EnsureApprovalStream(ctx)
}
// ConsumeApprovals 持久消费审批决定,每条交给 h 驱动 resume(队列组多副本安全)。
func (s *Subscriber) ConsumeApprovals(ctx context.Context, h func(context.Context, *contract.ApprovalDecision)) (func(context.Context), error) {
return s.inner.ConsumeApprovals(ctx, h)
}
// RequestModelConfig 向控制面(Gateway)取当前激活的对话模型配置。
func (s *Subscriber) RequestModelConfig(ctx context.Context) (*contract.ModelConfig, error) {
return s.inner.RequestConfig(ctx, contract.ConfigKindChat)
+55
View File
@@ -408,6 +408,61 @@ func (b *Bus) PublishApproval(dec *contract.ApprovalDecision) error {
return b.nc.Publish(contract.ApprovalSubject(dec.TaskID), data)
}
// EnsureApprovalStream 幂等地创建/更新审批决定流,持久捕获 sundynix.approval.>。
// 决定带 MaxAge 过期(审批不会拖超一天,避免流无限增长);让中断/恢复模型即便 dispatcher
// 在决定发出时离线,重连后仍能消费到决定续跑(core NATS 的 PublishApproval 也被本流捕获)。
func (b *Bus) EnsureApprovalStream(ctx context.Context) error {
_, err := b.js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
Name: contract.StreamApprovals,
Subjects: []string{contract.SubjectApprovalAll},
Storage: jetstream.FileStorage,
MaxAge: 24 * time.Hour,
})
return err
}
// ConsumeApprovals 在持久消费者上消费审批决定,队列组内负载均衡(多 dispatcher 副本下每条决定
// 只一个副本处理 resume)。每条决定派生独立 goroutine 处理(resume 续跑可能秒级),返回 stop 优雅停。
func (b *Bus) ConsumeApprovals(ctx context.Context, h func(context.Context, *contract.ApprovalDecision)) (drain func(context.Context), err error) {
cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamApprovals, jetstream.ConsumerConfig{
Durable: contract.ConsumerApprovals,
AckPolicy: jetstream.AckExplicitPolicy,
FilterSubject: contract.SubjectApprovalAll,
// resume 续跑可能再遇审批 / 出稿(LLM 秒级),AckWait 给足,未 ack 由重投兜底。
AckWait: 15 * time.Minute,
MaxAckPending: taskConcurrency(),
})
if err != nil {
return nil, fmt.Errorf("create approval consumer: %w", err)
}
var wg sync.WaitGroup
cc, err := cons.Consume(func(msg jetstream.Msg) {
var dec contract.ApprovalDecision
if json.Unmarshal(msg.Data(), &dec) != nil {
_ = msg.Term() // 脏数据,丢弃不重投
return
}
wg.Add(1)
go func() {
defer func() {
wg.Done()
if r := recover(); r != nil {
log.Printf("[bus] approval %s handler panic: %v", dec.TaskID, r)
_ = msg.Term()
}
}()
// ctx 派生自 Background:决定处理独立于触发它的请求/信号生命周期(同 ConsumeTasks)。
mctx := extractTrace(context.Background(), nats.Header(msg.Headers()))
h(mctx, &dec)
_ = msg.Ack()
}()
})
if err != nil {
return nil, fmt.Errorf("consume approvals: %w", err)
}
return func(dctx context.Context) { cc.Stop(); drainWait(&wg, dctx) }, nil
}
// WaitApproval 阻塞等待某任务的人工审批决定,直到收到、ctx 取消或超时。
// dispatcher 在审批节点调用:先订阅再等待(订阅早于决定到达,避免错过)。
// 超时返回 (nil, error) —— 调用方据安全默认(拒绝)处理。
+10 -2
View File
@@ -39,8 +39,16 @@ const (
SubjectTaskStatus = "sundynix.status.task"
// 人工审批(HITL)决定回传前缀:实际 sundynix.approval.<task_id>。
// core NATS pub-subUI 点批准/拒绝 → 网关发到此 → dispatcher 解除审批节点的阻塞。
SubjectApproval = "sundynix.approval"
// UI 点批准/拒绝 → 网关 nc.Publish 到此。阻塞模型经 core NATS WaitApproval 收;
// 中断/恢复模型经下方 JetStream 流持久捕获、由 dispatcher 持久消费者驱动 resume(抗离线)。
SubjectApproval = "sundynix.approval"
SubjectApprovalAll = "sundynix.approval.>" // 审批决定流捕获的通配
StreamApprovals = "SUNDYNIX_APPROVALS" // 审批决定 JetStream 流(持久,决定不因 dispatcher 离线而丢)
ConsumerApprovals = "approval-resumers" // 审批决定持久消费者(队列组:多副本下每条决定只一个副本处理 resume)
// BucketCheckpoints 是 HITL 持久化中断的 JetStream KV 桶名:存 compose 图 checkpoint
// (键=task_id)与 resume 记录(键=pending:task_id),dispatcher 重启后可据此恢复在途审批。
BucketCheckpoints = "SUNDYNIX_CHECKPOINTS"
// 自动化评测结果回写:dispatcher 评完经此广播,网关订阅落 PG 并供 UI 查询。core NATS pub-sub。
SubjectEval = "sundynix.eval.task"