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>
This commit is contained in:
Blizzard
2026-07-21 15:26:58 +08:00
parent ac38d5e663
commit 4500335da7
3 changed files with 272 additions and 71 deletions
+203
View File
@@ -0,0 +1,203 @@
package store
import (
"context"
"log"
"os"
"strings"
"time"
"gorm.io/gorm"
)
// 迁移机制(务实硬化,不引外部工具,见 production_readiness.md B1):
// 1. 整段迁移在 **PG advisory lock** 内串行 —— 多实例同时启动不再并发 ALTER/建索引竞争
// (此前无锁:两实例并发 AutoMigrate 一方报错即掉降级模式)。
// 2. 两处**破坏性** legacy 迁移(DROP TABLE CASCADE)默认**不在启动路径跑**,移到
// ALLOW_LEGACY_SCHEMA_MIGRATION=1 显式开关后;检测到旧 schema 但未开则只告警不动手。
// 3. **版本化 runner**AutoMigrate 之外的步骤(部分唯一索引、数据回填、将来 AutoMigrate
// 做不了的破坏性/数据迁移)登记在 schemaSteps,各跑一次并记入 schema_migration 表,
// 下次启动跳过。给了「有序、记录、跑一次」的真迁移语义,而不重写 gorm 结构体基线。
// migrationLockKey 是 pg_advisory_lock 的固定键(所有实例一致才能互斥)。
const migrationLockKey int64 = 20260721
// SchemaMigration 记录已应用的版本化迁移步骤(用模型而非裸 DDL,PG/sqlite 都可建,便于测试)。
type SchemaMigration struct {
ID int `gorm:"primaryKey"`
Name string `gorm:"size:128"`
AppliedAt time.Time `gorm:"autoCreateTime"`
}
func (SchemaMigration) TableName() string { return "sundynix_schema_migration" }
// migrationStep 是一步版本化迁移。fn 幂等更稳(存量库重跑无害),但 runner 靠 schema_migration
// 记录保证「已应用即跳过」,故不强求幂等——将来的破坏性步骤可以是非幂等的一次性 DDL。
type migrationStep struct {
id int
name string
fn func(*gorm.DB) error
}
// schemaSteps 是 AutoMigrate 之外的有序迁移。往后加破坏性/数据迁移 = 在末尾追加新 id,别改历史。
var schemaSteps = []migrationStep{
{1, "ledger_grant_ref_unique", func(db *gorm.DB) error {
// 支付入账幂等兜底闸:grant 分录按 ref(=订单号) 唯一。部分索引放行手工发放(ref 空)。
return db.Exec(`CREATE UNIQUE INDEX IF NOT EXISTS idx_ledger_grant_ref ON sundynix_credit_ledger (kind, ref) WHERE kind = 'grant' AND ref <> ''`).Error
}},
{2, "ledger_refund_ref_unique", func(db *gorm.DB) error {
// 退款幂等兜底闸:adjust 分录按 ref 唯一,与 grant 双闸对称。
return db.Exec(`CREATE UNIQUE INDEX IF NOT EXISTS idx_ledger_refund_ref ON sundynix_credit_ledger (kind, ref) WHERE kind = 'adjust' AND ref <> ''`).Error
}},
{3, "backfill_null_tenant_balance", func(db *gorm.DB) error {
// 回填历史 NULL 余额(credit_balance_micro 后加列,早于它的租户行为 NULL → 充值 NULL+N=NULL
// 永不到账)。存量重跑 WHERE IS NULL 无命中,安全。让「余额=SUM(ledger)」不变量重立。
return db.Exec(`UPDATE sundynix_tenant SET credit_balance_micro = COALESCE(
(SELECT SUM(credits_micro) FROM sundynix_credit_ledger l WHERE l.tenant_id = sundynix_tenant.id), 0)
WHERE credit_balance_micro IS NULL`).Error
}},
{4, "user_wechat_openid_unique", func(db *gorm.DB) error {
// 微信 openid 部分唯一索引:只约束非空。存量邮箱用户该列空串,普通唯一索引会互撞。
return db.Exec(`CREATE UNIQUE INDEX IF NOT EXISTS idx_user_wechat_openid ON sundynix_user (wechat_openid) WHERE wechat_openid <> ''`).Error
}},
}
// migratedModels 是 AutoMigrate 的基线模型清单(新增性 DDL,安全)。加表在此追加。
func migratedModels() []any {
return []any{
&User{}, &Task{}, &Eval{}, &LLMModel{}, &KB{}, &Doc{}, &Agent{}, &DocLink{}, &Pricing{}, &Prompt{},
&AuditLog{}, &GuardrailEvent{}, &Tenant{}, &TenantMember{}, &TenantInvite{}, &Space{}, &SpaceMember{},
&UsageEvent{}, &CreditLedger{}, &UsageRollup{}, &Setting{}, &CreditPack{}, &PaymentOrder{},
&RedeemCode{}, &SubscriptionPlan{}, &Subscription{}, &SchemaMigration{},
}
}
// runMigrations 在 advisory lock 内跑全部迁移:legacy(默认关) → AutoMigrate 基线 → 版本化步骤。
// 只有 AutoMigrate 失败才返回 error(→ 调用方降级);版本化步骤失败只记日志、下次启动重试。
func runMigrations(db *gorm.DB) error {
return withMigrationLock(db, func() error {
if allowLegacyMigration() {
migrateLegacyIntIDs(db)
migrateDocLinkToID(db)
} else {
warnIfLegacySchema(db)
}
if err := db.AutoMigrate(migratedModels()...); err != nil {
return err
}
if err := runVersionedMigrations(db, schemaSteps); err != nil {
log.Printf("[store] 版本化迁移失败: %v(下次启动重试)", err)
}
return nil
})
}
// withMigrationLock 取一条专用连接持 pg_advisory_lock,在锁内跑 fn,结束释放。
// 多实例同时启动只有一个进锁跑迁移,其余阻塞等待(避免并发 DDL 竞争)。
// 取锁给 60s 超时兜底:极端情况取不到就带告警继续(AutoMigrate/索引多为幂等,退一步不致命)。
func withMigrationLock(db *gorm.DB, fn func() error) error {
sqlDB, err := db.DB()
if err != nil {
return fn() // 拿不到底层连接(如测试用非标准驱动)→ 不阻塞,直接跑
}
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
defer cancel()
conn, err := sqlDB.Conn(ctx)
if err != nil {
log.Printf("[store] 取迁移锁连接失败,跳过加锁继续: %v", err)
return fn()
}
defer conn.Close()
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_lock($1)", migrationLockKey); err != nil {
log.Printf("[store] 取迁移 advisory lock 失败,跳过加锁继续: %v", err)
return fn()
}
defer func() {
uctx, ucancel := context.WithTimeout(context.Background(), 5*time.Second)
defer ucancel()
_, _ = conn.ExecContext(uctx, "SELECT pg_advisory_unlock($1)", migrationLockKey)
}()
return fn()
}
// runVersionedMigrations 顺序跑未应用的步骤,每步成功即记入 schema_migration 表。
// 某步失败即停(后续步骤可能依赖它),不记录 → 下次启动从该步重试。
func runVersionedMigrations(db *gorm.DB, steps []migrationStep) error {
if err := db.AutoMigrate(&SchemaMigration{}); err != nil {
return err
}
var ids []int
db.Model(&SchemaMigration{}).Pluck("id", &ids)
applied := make(map[int]bool, len(ids))
for _, id := range ids {
applied[id] = true
}
for _, s := range steps {
if applied[s.id] {
continue
}
if err := s.fn(db); err != nil {
return err
}
if err := db.Create(&SchemaMigration{ID: s.id, Name: s.name}).Error; err != nil {
return err
}
log.Printf("[store] 迁移 #%d(%s) 已应用", s.id, s.name)
}
return nil
}
// allowLegacyMigration 报告是否允许跑破坏性 legacy 迁移(默认关:不在启动路径 DROP 表)。
func allowLegacyMigration() bool {
v := os.Getenv("ALLOW_LEGACY_SCHEMA_MIGRATION")
return v == "1" || strings.EqualFold(v, "true")
}
// warnIfLegacySchema 检测到旧整型 id schema 但未开 legacy 开关时,只告警不动手(不 DROP)。
func warnIfLegacySchema(db *gorm.DB) {
var dt string
db.Raw(`SELECT data_type FROM information_schema.columns WHERE table_name='sundynix_model' AND column_name='id'`).Scan(&dt)
if dt == "bigint" || dt == "integer" {
log.Printf("[store] ⚠️ 检测到旧整型 id schema。破坏性迁移默认已关(不会自动 DROP 重建)。" +
"如确需迁移,设 ALLOW_LEGACY_SCHEMA_MIGRATION=1 后重启(会 DROP 并重建部分表)")
}
}
// migrateLegacyIntIDs 检测到旧整型 id 表则备份模型密钥、删旧表(AutoMigrate 随后按新规约重建)。
// **破坏性**:仅在 ALLOW_LEGACY_SCHEMA_MIGRATION=1 时经 runMigrations 调用。
func migrateLegacyIntIDs(db *gorm.DB) {
var dt string
db.Raw(`SELECT data_type FROM information_schema.columns WHERE table_name='sundynix_model' AND column_name='id'`).Scan(&dt)
if dt != "bigint" && dt != "integer" {
return // 全新库或已是新规约
}
log.Println("[store] 检测到旧整型 id 表,执行雪花 id 迁移(保模型密钥,重置其它测试表)")
var saved []map[string]any
db.Table("sundynix_model").Find(&saved)
for _, t := range []string{"sundynix_doc_link", "sundynix_doc", "sundynix_agent", "sundynix_kb", "sundynix_model", "sundynix_task", "sundynix_user"} {
db.Exec("DROP TABLE IF EXISTS " + t + " CASCADE")
}
_ = db.AutoMigrate(&LLMModel{}) // 先建模型表以回灌
for _, r := range saved {
s := func(k string) string { v, _ := r[k].(string); return v }
b, _ := r["active"].(bool)
_ = db.Create(&LLMModel{
Kind: s("kind"), Provider: s("provider"), BaseURL: s("base_url"),
APIKey: s("api_key"), Model: s("model"), Active: b,
}).Error
}
log.Printf("[store] 已回灌 %d 条模型配置(新雪花 id)", len(saved))
}
// migrateDocLinkToID 把旧的按名双链表迁到按 Doc.ID 关联的新表。**破坏性**(DROP 重建):
// 仅在 ALLOW_LEGACY_SCHEMA_MIGRATION=1 时调用。
func migrateDocLinkToID(db *gorm.DB) {
if !db.Migrator().HasTable("sundynix_doc_link") {
return
}
if db.Migrator().HasColumn(&DocLink{}, "from_id") {
return // 已是按 ID 关联的新 schema
}
log.Println("[store] 双链表升级为按文件 ID 关联,重建 sundynix_doc_link(链接随文档再入库重建)")
db.Exec("DROP TABLE IF EXISTS sundynix_doc_link CASCADE")
}
@@ -0,0 +1,65 @@
package store
import (
"errors"
"testing"
"gorm.io/gorm"
)
// 版本化 runner:首次跑全部、再跑全跳过、追加只跑新步。
func TestRunVersionedMigrations_RunsOnceThenSkips(t *testing.T) {
p := newTestStore(t)
var ran []int
steps := []migrationStep{
{1, "a", func(*gorm.DB) error { ran = append(ran, 1); return nil }},
{2, "b", func(*gorm.DB) error { ran = append(ran, 2); return nil }},
}
if err := runVersionedMigrations(p.db, steps); err != nil {
t.Fatal(err)
}
if len(ran) != 2 {
t.Fatalf("首次应跑 2 步,得 %v", ran)
}
ran = nil
if err := runVersionedMigrations(p.db, steps); err != nil {
t.Fatal(err)
}
if len(ran) != 0 {
t.Fatalf("第二次应全跳过(已记录),得 %v", ran)
}
steps = append(steps, migrationStep{3, "c", func(*gorm.DB) error { ran = append(ran, 3); return nil }})
if err := runVersionedMigrations(p.db, steps); err != nil {
t.Fatal(err)
}
if len(ran) != 1 || ran[0] != 3 {
t.Fatalf("应只跑新增步骤 #3,得 %v", ran)
}
}
// 某步失败即停:后续步骤不跑,且失败步骤不记录(下次从它重试);已成功步骤不重跑。
func TestRunVersionedMigrations_StopsOnFailure(t *testing.T) {
p := newTestStore(t)
var ran []int
steps := []migrationStep{
{1, "ok", func(*gorm.DB) error { ran = append(ran, 1); return nil }},
{2, "boom", func(*gorm.DB) error { return errors.New("boom") }},
{3, "after", func(*gorm.DB) error { ran = append(ran, 3); return nil }},
}
if err := runVersionedMigrations(p.db, steps); err == nil {
t.Fatal("步骤 #2 失败应返回 error")
}
if len(ran) != 1 {
t.Fatalf("步骤 #3 不该在 #2 失败后跑,得 %v", ran)
}
// 重试:#1 已记录跳过(不重复跑),#2 再次失败即停。
ran = nil
if err := runVersionedMigrations(p.db, steps); err == nil {
t.Fatal("重试仍应在 #2 失败")
}
if len(ran) != 0 {
t.Fatalf("已应用的 #1 不该重跑,得 %v", ran)
}
}
+4 -71
View File
@@ -61,84 +61,17 @@ func OpenPostgres(dsn string) *Postgres {
return &Postgres{}
}
tunePool(db) // 连接池上限,防高并发打爆 PG
// 一次性迁移:旧表用整型自增 id,与新雪花字符串 id 不兼容(AutoMigrate 不改主键类型)。
// 备份模型密钥(唯一不可再生的数据) → 重建全部表 → 回灌模型。其余为可重建的测试数据
migrateLegacyIntIDs(db)
migrateDocLinkToID(db)
if err := db.AutoMigrate(&User{}, &Task{}, &Eval{}, &LLMModel{}, &KB{}, &Doc{}, &Agent{}, &DocLink{}, &Pricing{}, &Prompt{}, &AuditLog{}, &GuardrailEvent{}, &Tenant{}, &TenantMember{}, &TenantInvite{}, &Space{}, &SpaceMember{}, &UsageEvent{}, &CreditLedger{}, &UsageRollup{}, &Setting{}, &CreditPack{}, &PaymentOrder{}, &RedeemCode{}, &SubscriptionPlan{}, &Subscription{}); err != nil {
log.Printf("[store] postgres AutoMigrate 失败,降级运行: %v", err)
// 迁移统一走 runMigrationsadvisory lock 内串行(多实例安全)→ AutoMigrate 基线 →
// 版本化步骤(索引/回填)。破坏性 legacy 迁移默认不跑(见 migrate.go
if err := runMigrations(db); err != nil {
log.Printf("[store] postgres 迁移失败,降级运行: %v", err)
return &Postgres{}
}
// 支付入账幂等兜底闸:grant 分录按 ref(=订单号) 唯一——支付回调是 at-least-once
// 订单状态机 CAS 是主闸,这里是第二道。部分索引:admin 手工发放 ref 为空、usage 分录不受影响。
if err := db.Exec(`CREATE UNIQUE INDEX IF NOT EXISTS idx_ledger_grant_ref ON sundynix_credit_ledger (kind, ref) WHERE kind = 'grant' AND ref <> ''`).Error; err != nil {
log.Printf("[store] 账本 grant/ref 唯一索引创建失败(重复入账兜底闸缺位): %v", err)
}
// 退款幂等兜底闸:adjust 分录带 ref(=订单号) 唯一——防重复退款冲销。
// 部分索引:admin 手工校正(GrantCredits 负数)ref 为空,不受约束;与 grant 双闸对称。
if err := db.Exec(`CREATE UNIQUE INDEX IF NOT EXISTS idx_ledger_refund_ref ON sundynix_credit_ledger (kind, ref) WHERE kind = 'adjust' AND ref <> ''`).Error; err != nil {
log.Printf("[store] 账本 adjust/ref 唯一索引创建失败(重复退款兜底闸缺位): %v", err)
}
// 回填历史 NULL 余额。credit_balance_micro 是后加的列,早于它创建的租户行值为 NULL,
// 而入账用的是 `余额 + N` —— SQL 里 NULL + N 仍是 NULL,于是这些租户**充值永远不到账**
// (分录照写、余额不动),且不报错。代码侧已改 coalesce 自愈,这里把存量一次修平,
// 让「余额 = SUM(ledger)」这条对账不变量重新成立。
if err := db.Exec(`UPDATE sundynix_tenant SET credit_balance_micro = COALESCE(
(SELECT SUM(credits_micro) FROM sundynix_credit_ledger l WHERE l.tenant_id = sundynix_tenant.id), 0)
WHERE credit_balance_micro IS NULL`).Error; err != nil {
log.Printf("[store] 历史 NULL 余额回填失败: %v", err)
}
// 微信 openid 部分唯一索引:只约束非空值。存量邮箱用户该列是空串 '' 而非 NULL
// 若建普通唯一索引,多个空串会互撞、AutoMigrate 直接失败(NULL 余额那次的同类坑)。
if err := db.Exec(`CREATE UNIQUE INDEX IF NOT EXISTS idx_user_wechat_openid ON sundynix_user (wechat_openid) WHERE wechat_openid <> ''`).Error; err != nil {
log.Printf("[store] 微信 openid 唯一索引创建失败: %v", err)
}
registerTenantScope(db) // 多租户:受租户模型的查询/创建自动按上下文注入 tenant_id(统一强制隔离)
log.Println("[store] postgres connected & migrated (雪花 id + 软删 规约)")
return &Postgres{db: db}
}
// migrateLegacyIntIDs 检测到旧整型 id 表则备份模型密钥、删旧表(AutoMigrate 随后按新规约重建)。
func migrateLegacyIntIDs(db *gorm.DB) {
var dt string
db.Raw(`SELECT data_type FROM information_schema.columns WHERE table_name='sundynix_model' AND column_name='id'`).Scan(&dt)
if dt != "bigint" && dt != "integer" {
return // 全新库或已是新规约
}
log.Println("[store] 检测到旧整型 id 表,执行雪花 id 迁移(保模型密钥,重置其它测试表)")
var saved []map[string]any
db.Table("sundynix_model").Find(&saved)
for _, t := range []string{"sundynix_doc_link", "sundynix_doc", "sundynix_agent", "sundynix_kb", "sundynix_model", "sundynix_task", "sundynix_user"} {
db.Exec("DROP TABLE IF EXISTS " + t + " CASCADE")
}
_ = db.AutoMigrate(&LLMModel{}) // 先建模型表以回灌
for _, r := range saved {
s := func(k string) string { v, _ := r[k].(string); return v }
b, _ := r["active"].(bool)
_ = db.Create(&LLMModel{
Kind: s("kind"), Provider: s("provider"), BaseURL: s("base_url"),
APIKey: s("api_key"), Model: s("model"), Active: b,
}).Error
}
log.Printf("[store] 已回灌 %d 条模型配置(新雪花 id)", len(saved))
}
// migrateDocLinkToID 把旧的按名双链表(from_name/to_name)迁到按 Doc.ID 关联的新表。
// 旧表无 from_id 列即判定为旧 schema:直接删表,由 AutoMigrate 重建;链接随文档再入库/编辑重建。
func migrateDocLinkToID(db *gorm.DB) {
if !db.Migrator().HasTable("sundynix_doc_link") {
return
}
if db.Migrator().HasColumn(&DocLink{}, "from_id") {
return // 已是按 ID 关联的新 schema
}
log.Println("[store] 双链表升级为按文件 ID 关联,重建 sundynix_doc_link(链接随文档再入库重建)")
db.Exec("DROP TABLE IF EXISTS sundynix_doc_link CASCADE")
}
// Enabled 报告是否处于真实持久化模式。
func (p *Postgres) Enabled() bool { return p.db != nil }