Files
sundynix-agentix/sundynix-gateway/cmd/server/main.go
T
Blizzard 940330cdb7 feat(eval): 评测结果回写升 JetStream 持久 —— 消灭最后一处 core NATS 回写 (P1)
完成度审计 P1 + 记忆 nats-durability 的既定规矩「计费/需落库的回写一律
JetStream+幂等,别fire-and-forget」。此前 eval 是全仓最后一处 core NATS pub-sub
回写:网关离线/慢消费者时评测结果直接丢——而质量趋势/门控都依赖它。

照抄已升级的 status/usage 范式(同为 dispatcher→gateway→PG 回写):
- contract 加 StreamEval/ConsumerEval;bus 加 EnsureEvalStream + ConsumeEval
  (durable consumer,AckExplicit,落库失败 Nak 重投自愈),PublishEval 改
  js.Publish 同步等 ack。删 core NATS 的 SubscribeEval。
- gateway/dispatcher 两个 wrapper 在 connect 时 EnsureEvalStream;gateway main
  的评测订阅从 SubscribeEval(fire-and-forget)换 ConsumeEval(handler 返 error→
  Nak),接入优雅停机 drain。
- 幂等前提已满足:SaveEval 按 task_id upsert,at-least-once 重投只覆盖不重复。

验证:e2e 去重测试(TestGatewayQueueDedup)升级到 ConsumeEval,50 条两副本合计
处理50次零重复;live 真端到端:提交任务→跑完→评测经新 eval 流落 PG(level=ok)
一次成功;两服务启动 eval 流 ensure 无报错。go build/vet/test 全绿。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-18 11:52:40 +08:00

197 lines
8.8 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/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 而非 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()
// 多租户回填:给存量用户补建默认租户 + 给受租户表存量行回填 tenant_id(幂等;每次启动跑一次)。
if n, err := db.BackfillDefaultTenants(context.Background()); err != nil {
log.Printf("[startup] 默认租户回填失败: %v", err)
} else if n > 0 {
log.Printf("[startup] 为 %d 个存量用户补建了默认租户", n)
}
if err := db.BackfillRowTenants(context.Background()); err != nil {
log.Printf("[startup] 存量行 tenant_id 回填失败: %v", err)
}
// 增量3 共享工作区:给每个租户成员补个人空间,再把存量 Agent re-key 到其个人空间(幂等;顺序:
// 空间先建、Agent 后迁——迁移内先回填 space_id 再建唯一索引,避免存量空 space_id 撞车)。
if n, err := db.BackfillPersonalSpaces(context.Background()); err != nil {
log.Printf("[startup] 个人空间回填失败: %v", err)
} else if n > 0 {
log.Printf("[startup] 为 %d 个租户成员补建了个人空间", n)
}
if err := db.MigrateAgentSpaces(context.Background()); err != nil {
log.Printf("[startup] Agent 空间作用域迁移失败: %v", err)
}
// KB/Doc/DocLink 的 PG 作用域迁移(space_id 回填 + 唯一索引换新)。存储层(Milvus/Bleve/Neo4j)
// 的存量向量另由一次性重灌迁移处理(不塞启动,避免每次重启重嵌 + 依赖 mcp-go 启动顺序)。
if err := db.MigrateKBSpaces(context.Background()); err != nil {
log.Printf("[startup] KB 空间作用域迁移失败: %v", err)
}
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 查询。
// JetStream at-least-once:落库失败 → Nak 重投自愈;UpdateTaskStatus 幂等(按 task_id 覆盖),重投无害。
statusDrain, serr := bus.ConsumeTaskStatus(context.Background(), func(ctx context.Context, ev *contract.TaskStatusEvent) error {
return db.UpdateTaskStatus(ctx, ev.TaskID, ev.Status, ev.Detail)
})
if serr != nil {
log.Printf("[gateway] consume task status: %v", serr)
}
// 评测闭环:持久消费 dispatcher 回写的自动化评测结果,落 PG 供 UI 查询 / 质量趋势 / 门控。
// JetStream at-least-once + SaveEval 按 task_id upsert 幂等;落库失败返 error → Nak 重投自愈。
evalDrain, everr := bus.ConsumeEval(context.Background(), func(ctx context.Context, ev *contract.EvalEvent) error {
flags, _ := json.Marshal(ev.Flags)
return db.SaveEval(ctx, &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,
})
})
if everr != nil {
log.Printf("[gateway] consume eval: %v", everr)
}
// 成本护栏:持久消费 dispatcher 回写的任务 token 用量(计费事实源,JetStream at-least-once + 幂等)。
usageDrain, uerr := bus.ConsumeUsage(context.Background(), func(ctx context.Context, ev *contract.UsageEvent) error {
if ev.UserID == "" || ev.TotalTok <= 0 {
return nil // 无效事件:Ack 丢弃,不重投
}
// 计费事实源:折算 credits + cost 幂等落持久明细(按 task_id 去重)。失败 → Nak 重投自愈。
inserted, err := db.SaveUsageEvent(ctx, ev)
if err != nil {
return err
}
// 仅"新插入"时累计 Redis 日计数(非幂等,重投不能重复累加;供提交前日预算门控)。
if inserted {
day := time.UnixMilli(ev.TS).Format("20060102")
if _, err := cache.AddUsage(ctx, ev.UserID, day, ev.TotalTok); err != nil {
log.Printf("[gateway] 累计用量 user=%s 失败: %v", ev.UserID, err) // best-effort,不影响已落库的计费
}
}
return nil
})
if uerr != nil {
log.Printf("[gateway] consume usage: %v", uerr)
}
// 入库工作队列:启动有界并发 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()
}
// 停状态/用量持久消费者(未 ack 的重启后由 JetStream 重投,幂等落库不丢不重)。
if statusDrain != nil {
statusDrain(context.Background())
}
if usageDrain != nil {
usageDrain(context.Background())
}
if evalDrain != nil {
evalDrain(context.Background())
}
log.Println("[gateway] 已优雅停机")
}
func envOr(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}