Files
sundynix-agentix/sundynix-gateway/internal/store/pgsql.go
T
Blizzard 4500335da7 refactor(store): 迁移路径务实硬化——advisory lock + 版本表 + 破坏性迁移移出启动路径(B1)
此前每次启动裸跑 AutoMigrate + 一串手写索引/回填,且启动路径上有两处 DROP TABLE
CASCADE 的 legacy 迁移;多实例并发启动无锁 → 并发 ALTER/建索引竞争,一方报错即掉降级。

务实硬化(不引外部工具,保留 gorm 结构体为源):
- PG advisory lock:整段迁移在 pg_advisory_lock 内串行,多实例同时启动只有一个进锁跑,
  其余阻塞等待。取锁 60s 超时兜底(取不到带告警继续,AutoMigrate/索引多幂等)。
- 破坏性 legacy 迁移移出默认路径:migrateLegacyIntIDs/migrateDocLinkToID(DROP TABLE
  CASCADE)默认不跑,仅 ALLOW_LEGACY_SCHEMA_MIGRATION=1 时执行;检测到旧 schema 但未开
  只告警不动手。现网早已是雪花 id 本就不触发,但从此不再是启动就可能 DROP。
- 版本化 runner:AutoMigrate 之外的步骤(3 个部分唯一索引 + NULL 余额回填)登记为
  schemaSteps,各跑一次并记入 sundynix_schema_migration 表,下次跳过;某步失败即停、
  不记录、下次重试。将来 AutoMigrate 做不了的破坏性/数据迁移在末尾追加新 id 即可。

pgsql.go 的迁移块收敛为一句 runMigrations(db)。带 runner 单测(跑一次/跳过/追加/失败即停)。
现网 DB 已有的索引/回填重跑无害(IF NOT EXISTS / WHERE IS NULL),跑后记录。build+vet+test 绿。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 15:26:58 +08:00

305 lines
11 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.
// Package store 封装 MainDB(PgSQL) 与 CacheDB(Redis) 的访问。
package store
import (
"context"
"errors"
"log"
"os"
"strconv"
"time"
"gorm.io/driver/postgres"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"gorm.io/gorm/schema"
"github.com/sundynix/sundynix-shared/contract"
)
// envInt 读正整数环境变量,缺省回退 def。
func envInt(key string, def int) int {
if v := os.Getenv(key); v != "" {
if n, err := strconv.Atoi(v); err == nil && n > 0 {
return n
}
}
return def
}
// tunePool 给连接池设上限:高并发下不至于无限开连接打爆 PG(max_connections 默认 100)。
// 各服务默认 25,可经 DB_MAX_OPEN_CONNS / DB_MAX_IDLE_CONNS 调整。
func tunePool(db *gorm.DB) {
if sqlDB, err := db.DB(); err == nil {
sqlDB.SetMaxOpenConns(envInt("DB_MAX_OPEN_CONNS", 25))
sqlDB.SetMaxIdleConns(envInt("DB_MAX_IDLE_CONNS", 5))
sqlDB.SetConnMaxLifetime(time.Hour)
}
}
// errStoreDisabled 表示 Postgres 处于降级(未连接)模式,写操作无法进行。
var errStoreDisabled = errors.New("postgres store disabled")
// Postgres 持有 MainDB 连接(Users / Billing / DSL)。
// db 为 nil 表示降级模式(连接失败时仍允许网关启动)。
type Postgres struct {
db *gorm.DB
}
// OpenPostgres 用 GORM 连接 MainDB 并自动迁移表结构。
// 表名统一 sundynix_ 前缀 + 单数(User→sundynix_user, Task→sundynix_task)。
// 连接失败不 fatal:返回降级实例,网关仍可启动(无 Docker 跑 demo 时即此路径)。
func OpenPostgres(dsn string) *Postgres {
db, err := gorm.Open(postgres.New(postgres.Config{DSN: dsn}), &gorm.Config{
NamingStrategy: schema.NamingStrategy{
TablePrefix: "sundynix_", // 所有表加前缀
SingularTable: true, // 单数表名
},
})
if err != nil {
log.Printf("[store] postgres 不可用,降级运行(不持久化): %v", err)
return &Postgres{}
}
tunePool(db) // 连接池上限,防高并发打爆 PG
// 迁移统一走 runMigrationsadvisory lock 内串行(多实例安全)→ AutoMigrate 基线 →
// 版本化步骤(索引/回填)。破坏性 legacy 迁移默认不跑(见 migrate.go)。
if err := runMigrations(db); err != nil {
log.Printf("[store] postgres 迁移失败,降级运行: %v", err)
return &Postgres{}
}
registerTenantScope(db) // 多租户:受租户模型的查询/创建自动按上下文注入 tenant_id(统一强制隔离)
log.Println("[store] postgres connected & migrated (雪花 id + 软删 规约)")
return &Postgres{db: db}
}
// Enabled 报告是否处于真实持久化模式。
func (p *Postgres) Enabled() bool { return p.db != nil }
// Ping 活性探测:底层连接池发一次 PingContext,验证 PG 此刻真的可达(非仅启动时连过)。
func (p *Postgres) Ping(ctx context.Context) bool {
if p.db == nil {
return false
}
sqlDB, err := p.db.DB()
if err != nil {
return false
}
return sqlDB.PingContext(ctx) == nil
}
// SaveTask 持久化一次任务提交(best-effort:降级模式下静默跳过)。
func (p *Postgres) SaveTask(ctx context.Context, owner, id, graph string) error {
if p.db == nil {
return nil
}
// TenantID 由 tenant 插件按请求 ctx 自动填;Owner 显式记录提交者(供个人工作台过滤)。
return p.db.WithContext(ctx).Create(&Task{Owner: owner, TaskID: id, Graph: graph, Status: contract.TaskSubmitted}).Error
}
// UpdateTaskStatus 流转任务状态(running/done/failed/timeout),由 dispatcher 经 NATS 回写驱动。
func (p *Postgres) UpdateTaskStatus(ctx context.Context, id, status, detail string) error {
if p.db == nil {
return nil
}
return p.db.WithContext(ctx).Model(&Task{}).
Where("task_id = ?", id).
Updates(map[string]any{"status": status, "detail": detail}).Error
}
// GetTaskStatus 取一条任务的当前状态(供 UI 轮询;不存在返回空串)。
func (p *Postgres) GetTaskStatus(ctx context.Context, id string) (status, detail string) {
if p.db == nil {
return "", ""
}
var t Task
if err := p.db.WithContext(ctx).Select("status", "detail").Where("task_id = ?", id).First(&t).Error; err != nil {
return "", ""
}
return t.Status, t.Detail
}
// SaveTaskOutput 收尾时落库最终模型输出(供历史复盘,best-effort)。
func (p *Postgres) SaveTaskOutput(ctx context.Context, id, output string) error {
if p.db == nil {
return nil
}
return p.db.WithContext(ctx).Model(&Task{}).Where("task_id = ?", id).Update("output", output).Error
}
// SaveTaskTrace 收尾时落库执行轨迹 JSON(供历史复盘,best-effort)。
func (p *Postgres) SaveTaskTrace(ctx context.Context, id, traceJSON string) error {
if p.db == nil {
return nil
}
return p.db.WithContext(ctx).Model(&Task{}).Where("task_id = ?", id).Update("trace", traceJSON).Error
}
// GetRunDetail 取一条任务持久化的输出 + 轨迹(历史复盘读库,不依赖 Redis TTL)。
func (p *Postgres) GetRunDetail(ctx context.Context, id string) (output, trace string) {
if p.db == nil {
return "", ""
}
var t Task
if err := p.db.WithContext(ctx).Select("output", "trace").Where("task_id = ?", id).First(&t).Error; err != nil {
return "", ""
}
return t.Output, t.Trace
}
// SaveEval 落库一条评测结果(按 task_id upsert:重评覆盖)。
func (p *Postgres) SaveEval(ctx context.Context, e *Eval) error {
if p.db == nil {
return nil
}
// 评测由 dispatcher 经 NATS 回写(background ctx,无请求租户)→ 从对应 task 复制 owner+tenant。
var t Task
if err := p.db.WithContext(ctx).Select("owner", "tenant_id").Where("task_id = ?", e.TaskID).First(&t).Error; err == nil {
e.Owner, e.TenantID = t.Owner, t.TenantID
}
return p.db.WithContext(ctx).Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "task_id"}},
DoUpdates: clause.AssignmentColumns([]string{
"owner", "tenant_id", "overall", "rule", "llm", "faithful", "level", "flags", "reason", "sources", "corrected", "updated_at",
}),
}).Create(e).Error
}
// GetEval 取一条任务的评测结果(不存在返回 nil)。
func (p *Postgres) GetEval(ctx context.Context, taskID string) *Eval {
if p.db == nil {
return nil
}
var e Eval
if err := p.db.WithContext(ctx).Where("task_id = ?", taskID).First(&e).Error; err != nil {
return nil
}
return &e
}
// CountTasks 返回已提交任务数(降级模式返回 0)。
func (p *Postgres) CountTasks(ctx context.Context) (int64, error) {
if p.db == nil {
return 0, nil
}
var n int64
err := p.db.WithContext(ctx).Model(&Task{}).Count(&n).Error
return n, err
}
// DayCount 是「某天 / 某状态 → 计数」的一行(工作台趋势/分布用)。
type DayCount struct {
Key string `json:"key"`
Count int64 `json:"count"`
}
// Overview 是工作台概览的聚合数据。任务/评测为实例级(Task 无 owner,单租户部署即全量),
// 知识库为 owner 级。降级模式(db==nil)返回零值。
type Overview struct {
TasksToday int64 `json:"tasks_today"`
TasksTotal int64 `json:"tasks_total"`
StatusCount []DayCount `json:"status_count"` // 近 7 天各终态分布
TaskTrend []DayCount `json:"task_trend"` // 近 7 天每日任务数(MM-DD)
EvalAvg float64 `json:"eval_avg"` // 综合分均值
FaithfulAvg float64 `json:"faithful_avg"` // 忠实度均值(仅有来源的)
EvalCount int64 `json:"eval_count"`
KBDocs int64 `json:"kb_docs"` // owner 文档数
KBCount int64 `json:"kb_count"` // owner 知识库数
}
// StatsOverview 聚合工作台概览(几条轻量查询)。owner 用于知识库口径。
func (p *Postgres) StatsOverview(ctx context.Context, owner string) *Overview {
o := &Overview{StatusCount: []DayCount{}, TaskTrend: []DayCount{}}
if p.db == nil {
return o
}
db := p.db.WithContext(ctx)
startOfDay := time.Now().Truncate(24 * time.Hour)
db.Model(&Task{}).Count(&o.TasksTotal)
db.Model(&Task{}).Where("created_at >= ?", startOfDay).Count(&o.TasksToday)
// 近 7 天每日任务数(按日期分组,缺的天补 0 在前端/此处处理)。
db.Model(&Task{}).
Select("to_char(created_at, 'MM-DD') as key, count(*) as count").
Where("created_at >= ?", time.Now().AddDate(0, 0, -6).Truncate(24*time.Hour)).
Group("key").Order("key").Scan(&o.TaskTrend)
// 近 7 天终态分布。
db.Model(&Task{}).
Select("status as key, count(*) as count").
Where("created_at >= ?", time.Now().AddDate(0, 0, -6)).
Group("status").Scan(&o.StatusCount)
// 评测均值(综合 + 忠实度仅算有来源的)。
var ev struct {
Avg float64
Faithful float64
N int64
}
db.Model(&Eval{}).Select("coalesce(avg(overall),0) as avg, coalesce(avg(nullif(faithful,0)),0) as faithful, count(*) as n").Scan(&ev)
o.EvalAvg, o.FaithfulAvg, o.EvalCount = ev.Avg, ev.Faithful, ev.N
if owner != "" {
db.Model(&Doc{}).Where("owner = ?", owner).Count(&o.KBDocs)
db.Model(&KB{}).Where("owner = ?", owner).Count(&o.KBCount)
}
return o
}
// RunRow 是「运行历史」一行:任务 + 其评测(LEFT JOIN,未评则 level 空)。
// 工作台「最近任务」与「运行」页共用它 —— 同一份数据只能有一个查法。
type RunRow struct {
TaskID string `json:"task_id"`
Status string `json:"status"`
Detail string `json:"detail"`
At time.Time `json:"at"`
EvalLevel string `json:"eval_level"`
EvalOverall float64 `json:"eval_overall"`
// Topic:报告类运行的主题。报告的 graph 是占位 DSL `{"topic":"…"}`,普通任务的 DSL
// 没有顶层 topic → 空串。否则运行历史里报告只能显示 report_<hex> 这种 id,读不出是啥。
Topic string `json:"topic"`
}
// RecentRuns 返回某用户最近 n 条运行(含评测分级,供「运行历史」列表)。
// 注:raw Table 查询绕过 gorm 模型回调 → 租户插件不生效,故此处**手动**按 owner(+ctx 租户) 过滤。
func (p *Postgres) RecentRuns(ctx context.Context, owner string, n int) []RunRow {
if p.db == nil {
return nil
}
var out []RunRow
q := p.db.WithContext(ctx).Table("sundynix_task as t").
Select("t.task_id, t.status, t.detail, t.created_at as at, "+
"coalesce(e.level,'') as eval_level, coalesce(e.overall,0) as eval_overall, "+
"coalesce(t.graph->>'topic','') as topic").
Joins("left join sundynix_eval e on e.task_id = t.task_id").
Where("t.deleted_at is null AND t.owner = ?", owner)
if tid := tenantFromCtx(ctx); tid != "" && !isSystemCtx(ctx) {
q = q.Where("t.tenant_id = ?", tid)
}
q.Order("t.created_at desc").Limit(n).Scan(&out)
return out
}
// SystemCounts 返回系统级计数(管理端 overview 用:全平台口径,非 owner 隔离)。
// 软删行由 GORM DeletedAt 作用域自动排除。降级模式返回零值。
func (p *Postgres) SystemCounts(ctx context.Context) (users, kbs, docs int64) {
if p.db == nil {
return 0, 0, 0
}
db := p.db.WithContext(ctx)
db.Model(&User{}).Count(&users)
db.Model(&KB{}).Count(&kbs)
db.Model(&Doc{}).Count(&docs)
return
}
// Close 释放底层连接。
func (p *Postgres) Close() {
if p.db == nil {
return
}
if sqlDB, err := p.db.DB(); err == nil {
_ = sqlDB.Close()
}
}