fix(v2): uq_delivery 索引显式升级迁移(AutoMigrate 不改同名索引)+ refund_id NOT NULL DEFAULT ''
AutoMigrate 只按名字判断索引是否存在,P2 时代库升级后同名 uq_delivery 仍停在 2 列,store.EnqueueDelivery 的 3 列 ON CONFLICT 每次调用(含普通 payment webhook)都报 "ON CONFLICT clause does not match any PRIMARY KEY or UNIQUE constraint"(reviewer 在 glebarez/sqlite 上实测复现)。新增 model.UpgradeWebhookDeliveryIndex,AutoMigrate 后显式检测/加宽索引 + backfill refund_id,main.go 与 OpenTestDB 两处调用点接入;RefundID 补 not null;default:'' 补上 NULL/"" 语义缺口。 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013nMthbVEmQquxBRKb9Fj8u
This commit is contained in:
@@ -0,0 +1,124 @@
|
||||
package model
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// UpgradeWebhookDeliveryIndex is an idempotent, dialect-aware schema-upgrade
|
||||
// step meant to run immediately after AutoMigrate (both in main.go's
|
||||
// autoMigrate and in OpenTestDB).
|
||||
//
|
||||
// Why this exists: GORM's AutoMigrate only checks an index's EXISTENCE by
|
||||
// name, never its column set. A P2-era database already has a 2-column
|
||||
// uq_delivery(out_trade_no,event_type) unique index; when the model widened
|
||||
// it to uq_delivery(out_trade_no,event_type,refund_id) (P4), AutoMigrate sees
|
||||
// "uq_delivery exists" and leaves the stale 2-column index untouched. From
|
||||
// then on store.EnqueueDelivery's 3-column `ON CONFLICT (out_trade_no,
|
||||
// event_type,refund_id)` matches no real unique constraint and fails on
|
||||
// EVERY call — including plain payment webhooks, not just refunds.
|
||||
//
|
||||
// This function explicitly detects that mismatch and repairs it:
|
||||
// 1. If uq_delivery already covers refund_id, no-op (safe to call every
|
||||
// startup / every test DB open).
|
||||
// 2. Otherwise: backfill any NULL refund_id to the empty string (pre-upgrade rows never
|
||||
// had the column, or had it added without a default), drop the stale
|
||||
// 2-column index, and recreate it to match the current 3-column model.
|
||||
func UpgradeWebhookDeliveryIndex(db *gorm.DB) error {
|
||||
covers, err := uqDeliveryCoversRefundID(db)
|
||||
if err != nil {
|
||||
return fmt.Errorf("model.UpgradeWebhookDeliveryIndex: detect: %w", err)
|
||||
}
|
||||
if covers {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Normalize before widening: rows created under the old 2-column index
|
||||
// may have NULL refund_id (payment events never set it, and a plain
|
||||
// AutoMigrate-added column can land NULL without an explicit default).
|
||||
// This also matches the RefundID `not null;default:''` column contract.
|
||||
if err := db.Exec(`UPDATE webhook_deliveries SET refund_id = '' WHERE refund_id IS NULL`).Error; err != nil {
|
||||
return fmt.Errorf("model.UpgradeWebhookDeliveryIndex: backfill refund_id: %w", err)
|
||||
}
|
||||
|
||||
if db.Migrator().HasIndex(&WebhookDelivery{}, "uq_delivery") {
|
||||
if err := db.Migrator().DropIndex(&WebhookDelivery{}, "uq_delivery"); err != nil {
|
||||
return fmt.Errorf("model.UpgradeWebhookDeliveryIndex: drop stale uq_delivery: %w", err)
|
||||
}
|
||||
}
|
||||
if err := db.Migrator().CreateIndex(&WebhookDelivery{}, "uq_delivery"); err != nil {
|
||||
return fmt.Errorf("model.UpgradeWebhookDeliveryIndex: create widened uq_delivery: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// uqDeliveryCoversRefundID reports whether the uq_delivery index on
|
||||
// webhook_deliveries already includes refund_id as one of its columns.
|
||||
// Dialect-aware: sqlite is queried via pragma_index_info (equivalent to
|
||||
// `PRAGMA index_info(uq_delivery)`, exercised by the sqlite branch of the
|
||||
// regression test); mysql via information_schema.statistics.
|
||||
//
|
||||
// If the index doesn't exist at all yet (fresh DB — AutoMigrate creates it
|
||||
// straight from the current 3-column model tags), this reports covers=true
|
||||
// so the caller treats it as nothing-to-widen; AutoMigrate already did the
|
||||
// right thing in that case.
|
||||
func uqDeliveryCoversRefundID(db *gorm.DB) (bool, error) {
|
||||
switch db.Dialector.Name() {
|
||||
case "sqlite":
|
||||
return sqliteUqDeliveryCoversRefundID(db)
|
||||
case "mysql":
|
||||
return mysqlUqDeliveryCoversRefundID(db)
|
||||
default:
|
||||
return false, fmt.Errorf("unsupported dialect %q for uq_delivery upgrade check", db.Dialector.Name())
|
||||
}
|
||||
}
|
||||
|
||||
// sqliteUqDeliveryCoversRefundID — tested branch (glebarez/sqlite, used in
|
||||
// all package tests via OpenTestDB).
|
||||
func sqliteUqDeliveryCoversRefundID(db *gorm.DB) (bool, error) {
|
||||
var cols []struct {
|
||||
Name string `gorm:"column:name"`
|
||||
}
|
||||
if err := db.Raw(`SELECT name FROM pragma_index_info(?)`, "uq_delivery").Scan(&cols).Error; err != nil {
|
||||
return false, err
|
||||
}
|
||||
if len(cols) == 0 {
|
||||
// Index doesn't exist yet — AutoMigrate will create it fresh from
|
||||
// the current (3-column) model; nothing for us to widen.
|
||||
return true, nil
|
||||
}
|
||||
for _, c := range cols {
|
||||
if c.Name == "refund_id" {
|
||||
return true, nil
|
||||
}
|
||||
}
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// mysqlUqDeliveryCoversRefundID — reasoned, not exercised by tests (no mysql
|
||||
// available in this environment; prod uses mysql, tests use glebarez/sqlite
|
||||
// per repo convention). information_schema.statistics lists one row per
|
||||
// (index_name, column) pair for the current database/table, which is the
|
||||
// mysql analogue of sqlite's pragma_index_info.
|
||||
func mysqlUqDeliveryCoversRefundID(db *gorm.DB) (bool, error) {
|
||||
var cols []struct {
|
||||
ColumnName string `gorm:"column:COLUMN_NAME"`
|
||||
}
|
||||
if err := db.Raw(
|
||||
`SELECT COLUMN_NAME FROM information_schema.statistics
|
||||
WHERE table_schema = DATABASE() AND table_name = ? AND index_name = ?`,
|
||||
"webhook_deliveries", "uq_delivery",
|
||||
).Scan(&cols).Error; err != nil {
|
||||
return false, err
|
||||
}
|
||||
if len(cols) == 0 {
|
||||
return true, nil
|
||||
}
|
||||
for _, c := range cols {
|
||||
if c.ColumnName == "refund_id" {
|
||||
return true, nil
|
||||
}
|
||||
}
|
||||
return false, nil
|
||||
}
|
||||
@@ -30,6 +30,9 @@ func OpenTestDB(t *testing.T) *gorm.DB {
|
||||
&Product{}, &ProductPrice{}); err != nil {
|
||||
t.Fatalf("migrate: %v", err)
|
||||
}
|
||||
if err := UpgradeWebhookDeliveryIndex(db); err != nil {
|
||||
t.Fatalf("upgrade webhook_deliveries uq_delivery index: %v", err)
|
||||
}
|
||||
sqlDB, _ := db.DB()
|
||||
t.Cleanup(func() { _ = sqlDB.Close() })
|
||||
return db
|
||||
|
||||
@@ -0,0 +1,122 @@
|
||||
package model_test
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
"github.com/glebarez/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
|
||||
"github.com/wangjia/pay/internal/model"
|
||||
"github.com/wangjia/pay/internal/store"
|
||||
)
|
||||
|
||||
// legacyWebhookDelivery 复刻 P2 时代 webhook_deliveries 的表结构:uq_delivery
|
||||
// 只覆盖 (out_trade_no,event_type) 两列,压根没有 refund_id 列。用它单独
|
||||
// AutoMigrate 出一张全新表,模拟"升级前"的生产库(当前 model.WebhookDelivery
|
||||
// 是 P4 才把 refund_id 并入 uq_delivery 的三列版本)。
|
||||
type legacyWebhookDelivery struct {
|
||||
model.Base
|
||||
OutTradeNo string `gorm:"size:64;not null;uniqueIndex:uq_delivery"`
|
||||
EventType string `gorm:"size:32;not null;uniqueIndex:uq_delivery"`
|
||||
BizSystem string `gorm:"index;size:32"`
|
||||
Payload string `gorm:"type:text"`
|
||||
Delivered bool `gorm:"index;default:false"`
|
||||
Attempts int
|
||||
LastError string `gorm:"size:255"`
|
||||
}
|
||||
|
||||
func (legacyWebhookDelivery) TableName() string { return "webhook_deliveries" }
|
||||
|
||||
var legacyDBCounter int64
|
||||
|
||||
// openLegacyDB 开一个全新的内存 sqlite 库,只灌 legacyWebhookDelivery 这张
|
||||
// "老版本"表结构 —— 不经过 model.OpenTestDB(它已经是当前 3 列 schema,没法
|
||||
// 用来复现"库还停在 2 列"的升级前状态)。
|
||||
func openLegacyDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
n := atomic.AddInt64(&legacyDBCounter, 1)
|
||||
dsn := fmt.Sprintf("file:legacytestdb_%d?mode=memory&cache=shared", n)
|
||||
db, err := gorm.Open(sqlite.Open(dsn),
|
||||
&gorm.Config{Logger: logger.Default.LogMode(logger.Silent), TranslateError: true})
|
||||
if err != nil {
|
||||
t.Fatalf("open legacy db: %v", err)
|
||||
}
|
||||
if err := db.AutoMigrate(&legacyWebhookDelivery{}); err != nil {
|
||||
t.Fatalf("migrate legacy schema: %v", err)
|
||||
}
|
||||
sqlDB, _ := db.DB()
|
||||
t.Cleanup(func() { _ = sqlDB.Close() })
|
||||
return db
|
||||
}
|
||||
|
||||
// TestUpgradeWidensLegacyUqDeliveryIndex 复现 reviewer 报的 Critical:GORM
|
||||
// AutoMigrate 只按名字判断索引是否存在,P2 时代库升级后同名 uq_delivery 还
|
||||
// 停在 2 列,store.EnqueueDelivery 的 3 列 ON CONFLICT 命中不到任何唯一约束,
|
||||
// 每次调用(含最普通的 payment webhook)都会报
|
||||
// "ON CONFLICT clause does not match any PRIMARY KEY or UNIQUE constraint"。
|
||||
//
|
||||
// fix 前跑本测试应在 (a) 处失败复现该错误(RED);fix 后
|
||||
// (model.UpgradeWebhookDeliveryIndex 显式加宽索引 + backfill refund_id)
|
||||
// 应全绿(GREEN)。
|
||||
func TestUpgradeWidensLegacyUqDeliveryIndex(t *testing.T) {
|
||||
db := openLegacyDB(t)
|
||||
|
||||
// 模拟一条"升级前"就存在的历史行:没有 refund_id 列,自然也没有值。
|
||||
if err := db.Exec(
|
||||
`INSERT INTO webhook_deliveries
|
||||
(out_trade_no, event_type, biz_system, payload, delivered, attempts, created_at, updated_at)
|
||||
VALUES ('PAY-LEGACY', 'payment.succeeded', 'pangolin', '{}', 0, 0, datetime('now'), datetime('now'))`,
|
||||
).Error; err != nil {
|
||||
t.Fatalf("seed legacy row: %v", err)
|
||||
}
|
||||
|
||||
// "升级":跑当前(3 列 refund_id)模型的 AutoMigrate,再跑显式 schema 升级步骤 ——
|
||||
// 与 main.go 的 autoMigrate() / model.OpenTestDB 同一顺序。
|
||||
if err := db.AutoMigrate(&model.WebhookDelivery{}); err != nil {
|
||||
t.Fatalf("automigrate current schema: %v", err)
|
||||
}
|
||||
if err := model.UpgradeWebhookDeliveryIndex(db); err != nil {
|
||||
t.Fatalf("UpgradeWebhookDeliveryIndex: %v", err)
|
||||
}
|
||||
|
||||
// (a) EnqueueDelivery 的 3 列 ON CONFLICT 必须能命中唯一约束:同订单同
|
||||
// event_type 下,1 条 payment(refund_id="")+ 2 条不同 refund_id 的退款
|
||||
// 事件 = 3 行(升级前的 bug 会让这三次调用全部报错)。
|
||||
ws := store.NewWebhookStore(db)
|
||||
if err := ws.EnqueueDelivery("PAY-1", "pangolin", "payment.succeeded", "", `{}`); err != nil {
|
||||
t.Fatalf("enqueue payment: %v", err)
|
||||
}
|
||||
if err := ws.EnqueueDelivery("PAY-1", "pangolin", "refund.succeeded", "rf-A", `{}`); err != nil {
|
||||
t.Fatalf("enqueue refund rf-A: %v", err)
|
||||
}
|
||||
if err := ws.EnqueueDelivery("PAY-1", "pangolin", "refund.succeeded", "rf-B", `{}`); err != nil {
|
||||
t.Fatalf("enqueue refund rf-B: %v", err)
|
||||
}
|
||||
list, err := ws.ListUndelivered(10)
|
||||
if err != nil {
|
||||
t.Fatalf("list: %v", err)
|
||||
}
|
||||
var forPay1 int
|
||||
for _, row := range list {
|
||||
if row.OutTradeNo == "PAY-1" {
|
||||
forPay1++
|
||||
}
|
||||
}
|
||||
// 期望 PAY-1 恰 3 行:1 payment + 2 refund(rf-A / rf-B);第 4 行是 PAY-LEGACY 的历史行。
|
||||
if forPay1 != 3 {
|
||||
t.Fatalf("want 3 undelivered rows for PAY-1 (1 payment + 2 distinct refund_id), got %d: %+v", forPay1, list)
|
||||
}
|
||||
|
||||
// (b) 历史行(升级前没有 refund_id 列)的 refund_id 应被 backfill 成 '',不是 NULL。
|
||||
var legacyRefundID string
|
||||
if err := db.Raw(`SELECT refund_id FROM webhook_deliveries WHERE out_trade_no = 'PAY-LEGACY'`).
|
||||
Scan(&legacyRefundID).Error; err != nil {
|
||||
t.Fatalf("read legacy row: %v", err)
|
||||
}
|
||||
if legacyRefundID != "" {
|
||||
t.Fatalf("legacy row refund_id want '', got %q", legacyRefundID)
|
||||
}
|
||||
}
|
||||
@@ -6,7 +6,7 @@ type WebhookDelivery struct {
|
||||
Base
|
||||
OutTradeNo string `gorm:"size:64;not null;uniqueIndex:uq_delivery" json:"out_trade_no"`
|
||||
EventType string `gorm:"size:32;not null;uniqueIndex:uq_delivery" json:"event_type"`
|
||||
RefundID string `gorm:"size:64;uniqueIndex:uq_delivery" json:"refund_id,omitempty"` // 退款事件的幂等维度;payment 事件为空
|
||||
RefundID string `gorm:"size:64;not null;default:'';uniqueIndex:uq_delivery" json:"refund_id,omitempty"` // 退款事件的幂等维度;payment 事件为空
|
||||
BizSystem string `gorm:"index;size:32" json:"biz_system"`
|
||||
Payload string `gorm:"type:text" json:"payload"` // 已序列化的领域 JSON(含 event_type)
|
||||
Delivered bool `gorm:"index;default:false" json:"delivered"`
|
||||
|
||||
@@ -110,6 +110,11 @@ func autoMigrate(db *gorm.DB) {
|
||||
); err != nil {
|
||||
log.Fatalf("自动迁移失败: %v", err)
|
||||
}
|
||||
// AutoMigrate 只按名字判断索引是否存在,不会把老库同名的 2 列 uq_delivery
|
||||
// 加宽到当前模型的 3 列(out_trade_no,event_type,refund_id);显式补一刀。
|
||||
if err := model.UpgradeWebhookDeliveryIndex(db); err != nil {
|
||||
log.Fatalf("升级 uq_delivery 索引失败: %v", err)
|
||||
}
|
||||
log.Println("AutoMigrate 完成")
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user