Files
sundynix-agentix/sundynix-gateway/cmd/server/main.go
T
Blizzard ecf4a80466 feat(llm): T2.1 模型路由 + Fallback —— 单 provider 抖动不再整体宕
现状:单 provider,一家 API 抖动/挂掉全平台不可用。加主备 failover:主模型调用失败
自动按序切备用,compose/ReAct/Chat 全路径透明白嫖。

- llm/failover.go: failoverModel 把多个 ToolCallingChatModel 串成主备链,按序调用、
  遇错切下一个;它本身是 model.ToolCallingChatModel 故全路径透明。调用方主动取消
  (ctx.Err()!=nil) 不切;模型自身请求超时走内部 ctx、不污染父 ctx 故仍 failover。
  局限(v1):Stream 仅建流同步报错时切(已回流 token 的中途失败不切)。
- llm/pool.go: SetConfig 用激活配置(含 Fallbacks)重建——主+可用备用串成 failover 链,
  无备用则直接用主;备用单个构建失败跳过不影响主链。
- contract.ModelConfig: 加 Fallbacks 字段(骑在主配置里下发,不改任何 bus/ServeConfig 签名)。
- gateway store.ActiveConfig: chat 把"其它已登记 chat 模型"按序填进 Fallbacks;
  provide(main) + broadcast(admin) 共用 → 注册多个 chat 模型即自动成主备。
- bus.decryptConfig: 备用模型的 api_key(密文)一并解密。

测试:failover 单测(主可用不调备/主挂切备/全挂报错/取消不切/Stream 切备/WithTools 链)。
live 验证:active=死 ollama 主 + deepseek 备 → 任务连主拒连→自动切 deepseek→4s 完成。
DEPTH_ROADMAP T2.1(admin 注册多模型即主备,无需新 UI)。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-30 09:26:18 +08:00

130 lines
5.4 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 server 启动 sundynix-gateway —— 第 2 层业务网关 / 统一接入层。
package main
import (
"context"
"encoding/json"
"log"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/sundynix/sundynix-gateway/internal/blob"
"github.com/sundynix/sundynix-gateway/internal/nats"
"github.com/sundynix/sundynix-gateway/internal/router"
"github.com/sundynix/sundynix-gateway/internal/store"
sharedbus "github.com/sundynix/sundynix-shared/bus"
"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-gateway") // 结构化日志 + 链路感知(日志带 trace_id)
// 链路追踪:HTTP 入口 + 跨 NATS 传播的根。退出前 flush 残留 span。
shutdownTrace, _ := otelx.Init(context.Background(), "sundynix-gateway")
defer func() { _ = shutdownTrace(context.Background()) }()
natsURL := envOr("NATS_URL", "nats://localhost:4222")
// DSN 用 127.0.0.1 而非 localhostpgx 每新建连接解析一次主机名,高并发下连接池扩容时
// 大量 localhost DNS 查询会在请求超时里被取消而雪崩(见 LOAD_TEST_REPORT §4.3)。IP 免 DNS。
pgDSN := envOr("POSTGRES_DSN", "postgres://sundynix:sundynix@127.0.0.1:5432/sundynix?sslmode=disable")
redisAddr := envOr("REDIS_ADDR", "localhost:6379")
// 对象存储(大文档正文):默认连 docker-compose 暴露的 MinIO;连不上则降级内联存 PG。
blobStore := blob.Open(
envOr("MINIO_ENDPOINT", "localhost:9000"),
envOr("MINIO_ACCESS_KEY", "minioadmin"),
envOr("MINIO_SECRET_KEY", "minioadmin"),
envOr("MINIO_BUCKET", "sundynix-docs"),
)
db := store.OpenPostgres(pgDSN) // MainDB: Users / Billing / DSL(连不上则降级)
defer db.Close()
cache := store.OpenRedis(redisAddr) // CacheDB: Session / Rate Limit(连不上则降级)
defer cache.Close()
bus := nats.MustConnect(natsURL) // 接入 NATS 零拷贝骨干网 + 声明任务流
defer bus.Close()
// 配置控制面:按 kind 响应消费方(Dispatcher=chat / mcp-go=embedding)的配置请求。
for _, kind := range []string{contract.ConfigKindChat, contract.ConfigKindEmbedding} {
k := kind
if _, err := bus.ServeConfig(k, func() *contract.ModelConfig {
return db.ActiveConfig(context.Background(), k) // chat 含 Fallbacks(其它模型作备用)
}); err != nil {
log.Printf("[gateway] serve %s config: %v", k, err)
}
}
// 任务生命周期:订阅 dispatcher 回写的状态流转(running/done/failed/timeout),落 PG 供 UI 查询。
if _, err := bus.SubscribeTaskStatus(func(ev *contract.TaskStatusEvent) {
if err := db.UpdateTaskStatus(context.Background(), ev.TaskID, ev.Status, ev.Detail); err != nil {
log.Printf("[gateway] 更新任务状态 %s=%s 失败: %v", ev.TaskID, ev.Status, err)
}
}); err != nil {
log.Printf("[gateway] subscribe task status: %v", err)
}
// 评测闭环:订阅 dispatcher 回写的自动化评测结果,落 PG 供 UI 查询 / 质量趋势 / 门控。
if _, err := bus.SubscribeEval(func(ev *contract.EvalEvent) {
flags, _ := json.Marshal(ev.Flags)
if err := db.SaveEval(context.Background(), &store.Eval{
TaskID: ev.TaskID, Overall: ev.Overall, Rule: ev.Rule, LLM: ev.LLM, Faithful: ev.Faithful,
Level: ev.Level, Flags: string(flags), Reason: ev.Reason, Sources: ev.Sources, Corrected: ev.Corrected,
}); err != nil {
log.Printf("[gateway] 落库评测 %s 失败: %v", ev.TaskID, err)
}
}); err != nil {
log.Printf("[gateway] subscribe eval: %v", err)
}
// 成本护栏:订阅 dispatcher 回写的任务 token 用量,按用户按天累计到 Redis(供提交前日预算门控 / 计费)。
if _, err := bus.SubscribeUsage(func(ev *contract.UsageEvent) {
if ev.UserID == "" || ev.TotalTok <= 0 {
return
}
day := time.UnixMilli(ev.TS).Format("20060102")
if _, err := cache.AddUsage(context.Background(), ev.UserID, day, ev.TotalTok); err != nil {
log.Printf("[gateway] 累计用量 user=%s 失败: %v", ev.UserID, err)
}
}); err != nil {
log.Printf("[gateway] subscribe usage: %v", err)
}
r := router.New(db, cache, bus, blobStore)
addr := envOr("GATEWAY_ADDR", ":8080")
srv := &http.Server{Addr: addr, Handler: r}
// 后台监听;ListenAndServe 在 Shutdown 后返回 ErrServerClosed(正常退出)。
go func() {
log.Printf("[gateway] listening on %s", addr)
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Fatalf("[gateway] listen: %v", err)
}
}()
// 优雅停机:收到信号 → 停止收新连接、drain 在途请求(至多 DrainTimeout),再退出
// defer 的 db/cache/bus/blob Close 随后执行,确保 HTTP 排空后才断后端连接)。
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
<-ctx.Done()
to := sharedbus.DrainTimeout()
log.Printf("[gateway] 收到停机信号,优雅关闭 HTTP(drain 在途请求 ≤%s)…", to)
sctx, cancel := context.WithTimeout(context.Background(), to)
defer cancel()
if err := srv.Shutdown(sctx); err != nil {
log.Printf("[gateway] HTTP 优雅关闭超时/出错: %v", err)
}
log.Println("[gateway] 已优雅停机")
}
func envOr(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}