Files
sundynix-agentix/sundynix-dispatcher/cmd/dispatcher/main.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

117 lines
4.9 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.
// Command dispatcher 启动 sundynix-dispatcher —— 第 4 层 AI Agent 调度集群。
package main
import (
"context"
"encoding/json"
"log"
"os"
"os/signal"
"syscall"
"time"
"github.com/sundynix/sundynix-dispatcher/internal/eino"
"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"
)
func main() {
secrets.MustHaveKeyInProd() // 生产须配 SUNDYNIX_SECRET_KEY 以解密下发的 api_key 密文
otelx.SetupSlog("sundynix-dispatcher") // 结构化日志 + 链路感知(任务日志带 trace_id)
// 链路追踪:消费 span 续上 gateway 的 trace,再向下展开节点/工具/LLM span。
shutdownTrace, _ := otelx.Init(context.Background(), "sundynix-dispatcher")
defer func() { _ = shutdownTrace(context.Background()) }()
natsURL := envOr("NATS_URL", "nats://localhost:4222")
pool := llm.NewPool() // LLM Pool: vLLM / Ollama 集群
breaker := harness.NewCircuitBreaker() // Harness: 熔断降级中心
// Harness: LLM 自动化评测(规则 + LLM-as-judge,模型就绪时启用)。
llmChat := func(ctx context.Context, sys, user string) (string, error) {
return pool.Chat(ctx, []llm.ChatMessage{{Role: "system", Content: sys}, {Role: "user", Content: user}})
}
eval := harness.NewEvaluator(pool.Ready, llmChat)
// Harness: 输入护栏 Tier2 —— 对网关标记的灰区输入做 LLM 越狱裁决(模型就绪时启用)。
guardian := harness.NewClassifier(pool.Ready, llmChat)
sub := dnats.MustConnect(natsURL)
defer sub.Close()
// 配置控制面:先订阅热更新,再后台重试拉初始配置(容忍 gateway 晚于本服务启动,
// 避免一次性请求扑空后只能干等热更新 → 降级桩跑全程)。
if _, err := sub.SubscribeModelConfigUpdated(pool.SetConfig); err != nil {
log.Printf("[dispatcher] subscribe model config: %v", err)
}
go sub.FetchModelConfigWithRetry(context.Background(), pool.SetConfig)
// sub 同时作为 Token 回流(TokenSink)、MCP 工具调用(ToolCaller)、执行事件(ExecSink)、
// 任务状态回写(StatusSink)、HITL 审批等待(ApprovalWaiter)与评测落库(EvalSink)出口。
orch, err := eino.NewOrchestrator(pool, breaker, eval, sub, sub, sub, sub, sub, sub)
if err != nil {
log.Fatalf("[dispatcher] build eino graph: %v", err)
}
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 {
data, _ := json.Marshal(map[string]any{
"ok": true,
"service": "dispatcher",
"model": pool.ModelName(),
"ready": pool.Ready(),
"uptime_s": int(time.Since(startedAt).Seconds()),
})
return data
}); herr != nil {
log.Printf("[dispatcher] serve health: %v", herr)
} else {
defer func() { _ = unsub() }()
}
// 监听退出信号,优雅停止消费。
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
log.Println("[dispatcher] consuming sundynix.tasks.* (Ctrl-C to quit)")
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 {
if v := os.Getenv(key); v != "" {
return v
}
return def
}