From 00e032f577015c68bda54fda03243c386e82c900 Mon Sep 17 00:00:00 2001 From: wangjia <809946525@qq.com> Date: Fri, 10 Jul 2026 16:42:29 +0800 Subject: [PATCH] =?UTF-8?q?fix(v2):=20uq=5Fdelivery=20=E7=B4=A2=E5=BC=95?= =?UTF-8?q?=E6=98=BE=E5=BC=8F=E5=8D=87=E7=BA=A7=E8=BF=81=E7=A7=BB(AutoMigr?= =?UTF-8?q?ate=20=E4=B8=8D=E6=94=B9=E5=90=8C=E5=90=8D=E7=B4=A2=E5=BC=95)+?= =?UTF-8?q?=20refund=5Fid=20NOT=20NULL=20DEFAULT=20''?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 Claude-Session: https://claude.ai/code/session_013nMthbVEmQquxBRKb9Fj8u --- internal/model/schema_upgrade.go | 124 +++++++++++++++++++++++++++++ internal/model/testdb.go | 3 + internal/model/upgrade_test.go | 122 ++++++++++++++++++++++++++++ internal/model/webhook_delivery.go | 2 +- main.go | 5 ++ 5 files changed, 255 insertions(+), 1 deletion(-) create mode 100644 internal/model/schema_upgrade.go create mode 100644 internal/model/upgrade_test.go diff --git a/internal/model/schema_upgrade.go b/internal/model/schema_upgrade.go new file mode 100644 index 0000000..9e1152f --- /dev/null +++ b/internal/model/schema_upgrade.go @@ -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 +} diff --git a/internal/model/testdb.go b/internal/model/testdb.go index 6fa61c5..1aa1569 100644 --- a/internal/model/testdb.go +++ b/internal/model/testdb.go @@ -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 diff --git a/internal/model/upgrade_test.go b/internal/model/upgrade_test.go new file mode 100644 index 0000000..6e9896f --- /dev/null +++ b/internal/model/upgrade_test.go @@ -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) + } +} diff --git a/internal/model/webhook_delivery.go b/internal/model/webhook_delivery.go index f313bd6..56ebe91 100644 --- a/internal/model/webhook_delivery.go +++ b/internal/model/webhook_delivery.go @@ -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"` diff --git a/main.go b/main.go index 4519b17..bd4246b 100644 --- a/main.go +++ b/main.go @@ -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 完成") }