ac38d5e663
部署前生产级审计(可靠性/数据层/安全三路)后,清掉 7 处代码级硬伤: A1 后台定时器 goroutine 无 panic recover → 单个 DB panic 崩整个 gateway。加 safeGo/ safeCall,包住订阅/掉单补偿/微信推送/探针 goroutine,单轮 tick 再兜一层。 A2 提示词控制面(建/激活/停用,热广播全服务)只 RequireAuth → 任意登录用户改全局提示词。 三写端点+列表挂 RequireAdmin。 A3 HITL 审批端点无角色门 → viewer 可放行烧钱执行。加 RequireTenantRole(member)。 A4 审计/护栏列表 limit 无校验,limit=-1 让 gorm 取消 LIMIT 全表扫。加 clampLimit/ clampOffset,AdminTasks/AdminSpaces 补上界。 A5 限流 Redis 一挂就完全放行(fail-open)。加进程内固定窗口兜底(fail-safe) + 登录/注册 按 IP 专用严限流(10/min)。 A6 公开 by-id 端点(stream/exec/report导出/kb导入流)无鉴权无租户过滤。加 AuthFromHeaderOrQuery(从 ?token= 取 JWT) + task/report 按 owner 归属校验;桌面端 5 处 EventSource/下载 URL 经 tokenQuery 附 JWT。 A7 文件上传无大小上限(整文件进内存 OOM 面) → 50MB 闸(KB_MAX_UPLOAD_BYTES)+ LimitReader; http.Server 加 ReadHeaderTimeout/ReadTimeout/MaxHeaderBytes(不设 WriteTimeout 保 SSE)。 带单测:clampLimit/safeCall/procLimiter/AuthFromHeaderOrQuery/TaskOwner。 build+vet+全量 test 绿;desktop tsc 绿。B(迁移工具/实时探针/出网韧性/登录锁定/leader选举) 与 C(TLS/PG HA/K8s/备份自动化/可观测)分期后做,参照 production_readiness.md。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
205 lines
9.2 KiB
Go
205 lines
9.2 KiB
Go
// Command server 启动 sundynix-gateway —— 第 2 层业务网关 / 统一接入层。
|
||
package main
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"log"
|
||
"net/http"
|
||
"os"
|
||
"os/signal"
|
||
"syscall"
|
||
"time"
|
||
|
||
"github.com/sundynix/sundynix-shared/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()
|
||
// 多租户回填:给存量用户补建默认租户 + 给受租户表存量行回填 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")
|
||
// 慢读/Slowloris 防护:限制读头/读体时间与头大小。**不设 WriteTimeout**——会掐断
|
||
// SSE 长连(/tasks/:id/stream、/exec)。ReadTimeout 取 60s 容纳文件上传体。
|
||
srv := &http.Server{
|
||
Addr: addr,
|
||
Handler: r,
|
||
ReadHeaderTimeout: 10 * time.Second,
|
||
ReadTimeout: 60 * time.Second,
|
||
MaxHeaderBytes: 1 << 20, // 1MB
|
||
}
|
||
|
||
// 后台监听;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
|
||
}
|