diff --git a/sundynix-dispatcher/cmd/dispatcher/main.go b/sundynix-dispatcher/cmd/dispatcher/main.go index 42db8aa..d5dc5ff 100644 --- a/sundynix-dispatcher/cmd/dispatcher/main.go +++ b/sundynix-dispatcher/cmd/dispatcher/main.go @@ -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 { diff --git a/sundynix-dispatcher/internal/eino/compose_compiler.go b/sundynix-dispatcher/internal/eino/compose_compiler.go index 0613be4..e33ec53 100644 --- a/sundynix-dispatcher/internal/eino/compose_compiler.go +++ b/sundynix-dispatcher/internal/eino/compose_compiler.go @@ -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) + }() +} diff --git a/sundynix-dispatcher/internal/eino/compose_compiler_test.go b/sundynix-dispatcher/internal/eino/compose_compiler_test.go index 9a5035e..9a5e22a 100644 --- a/sundynix-dispatcher/internal/eino/compose_compiler_test.go +++ b/sundynix-dispatcher/internal/eino/compose_compiler_test.go @@ -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("批准后应续跑下游 agent,sink=%q", fs.text()) + } + if st.last() != contract.TaskDone { + t.Fatalf("续跑成功应判 done,got %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) { diff --git a/sundynix-dispatcher/internal/nats/subscriber.go b/sundynix-dispatcher/internal/nats/subscriber.go index 6bf0c5d..551c6cf 100644 --- a/sundynix-dispatcher/internal/nats/subscriber.go +++ b/sundynix-dispatcher/internal/nats/subscriber.go @@ -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.CheckpointKV(Get/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) diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index ba8014d..ddd5251 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -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) —— 调用方据安全默认(拒绝)处理。 diff --git a/sundynix-shared/contract/task.go b/sundynix-shared/contract/task.go index 276570a..cbd0b0c 100644 --- a/sundynix-shared/contract/task.go +++ b/sundynix-shared/contract/task.go @@ -39,8 +39,16 @@ const ( SubjectTaskStatus = "sundynix.status.task" // 人工审批(HITL)决定回传前缀:实际 sundynix.approval.。 - // core NATS pub-sub:UI 点批准/拒绝 → 网关发到此 → 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"