// 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/handler" "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 而非 localhost:pgx 每新建连接解析一次主机名,高并发下连接池扩容时 // 大量 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) } } // Prompt 控制面:响应各服务「取激活 prompt」请求(DB 激活集;空则服务用内置默认)。 if _, err := bus.ServePrompts(func() map[string]string { return db.ActivePrompts(context.Background()) }); err != nil { log.Printf("[gateway] serve prompts: %v", 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) } // 入库工作队列:启动有界并发 worker 池消费 JetStream 入库作业(背压防 OOM、崩溃重投兜底)。 ingestDrain, err := handler.StartIngestWorkers(context.Background(), db, cache, bus, blobStore) if err != nil { log.Printf("[gateway] 启动入库 worker 失败(入库不可用): %v", err) } else { log.Printf("[gateway] 入库 worker 池就绪(JetStream 持久队列)") } 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) } // HTTP 排空后再 drain 在途入库作业(≤to);超时未完的不丢——AckWait 30min 内重启会重投续跑。 if ingestDrain != nil { log.Printf("[gateway] drain 在途入库作业 ≤%s(超时未完者重启后由队列重投)…", to) ictx, icancel := context.WithTimeout(context.Background(), to) ingestDrain(ictx) icancel() } log.Println("[gateway] 已优雅停机") } func envOr(key, def string) string { if v := os.Getenv(key); v != "" { return v } return def }