feat: 增加规格匹配建议与决策审计 (#90)

This commit is contained in:
chengma
2026-08-10 12:41:17 +08:00
parent 5abe81e128
commit 5a77bde534
15 changed files with 680 additions and 46 deletions
+12
View File
@@ -50,3 +50,15 @@ func UpsertSpecMapping(q Execer, m model.SpecMapping) error {
}
return nil
}
func InsertSpecMappingDecision(q Execer, d model.SpecMappingDecision) error {
_, err := q.Exec(`INSERT INTO spec_mapping_decisions
(shopee_goods_id,spec_key,pdd_goods_id,rules_version,suggested_option_key,
chosen_option_key,accepted,decided_by,decided_at) VALUES (?,?,?,?,?,?,?,?,?)`,
d.ShopeeGoodsID, d.SpecKey, d.PddGoodsID, d.RulesVersion, nullableText(d.SuggestedOptionKey),
d.ChosenOptionKey, d.Accepted, nullableText(d.DecidedBy), d.DecidedAt)
if err != nil {
return fmt.Errorf("写入规格匹配决策审计失败: %w", err)
}
return nil
}
+71 -3
View File
@@ -19,7 +19,7 @@ import (
"cmautobuy/admin/spec"
)
const mysqlSchemaVersion = 3
const mysqlSchemaVersion = 4
// OpenMySQL 打开生产 MySQL 8 数据库。错误信息绝不包含完整 DSN 或密码。
func OpenMySQL(cfg config.DatabaseConfig) (*sql.DB, error) {
@@ -387,6 +387,21 @@ const mysqlSchemaV3SpecMappings = `CREATE TABLE IF NOT EXISTS spec_mappings (
KEY idx_spec_mappings_pdd (pdd_goods_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`
const mysqlSchemaV4Decisions = `CREATE TABLE IF NOT EXISTS spec_mapping_decisions (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
shopee_goods_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
spec_key VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
pdd_goods_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
rules_version VARCHAR(32) NOT NULL,
suggested_option_key VARCHAR(191) COLLATE utf8mb4_bin,
chosen_option_key VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
accepted TINYINT NOT NULL,
decided_by VARCHAR(191),
decided_at VARCHAR(35) NOT NULL,
CONSTRAINT chk_spec_mapping_decisions_accepted CHECK (accepted IN (0,1)),
KEY idx_decisions_mapping (shopee_goods_id, spec_key, pdd_goods_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`
// MigrateMySQL 建立或升级 MySQL schema。生产迁移只能在这里追加新版本。
func MigrateMySQL(db *sql.DB) error {
if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS schema_migrations (
@@ -437,17 +452,37 @@ func MigrateMySQL(db *sql.DB) error {
if err := migrateMySQLV3(db); err != nil {
return fmt.Errorf("执行 MySQL schema v3 失败: %w", err)
}
if err := CheckMySQLSchema(db); err != nil {
if err := checkMySQLSchemaV3(db); err != nil {
return fmt.Errorf("MySQL schema v3 自检失败,未记录版本: %w", err)
}
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`,
3, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
return fmt.Errorf("记录 MySQL schema v3 失败: %w", err)
}
current = 3
}
if current < 4 {
if _, err := db.Exec(mysqlSchemaV4Decisions); err != nil {
return fmt.Errorf("执行 MySQL schema v4 失败: %w", err)
}
if err := checkMySQLV4Shape(db); err != nil {
return fmt.Errorf("MySQL schema v4 自检失败,未记录版本: %w", err)
}
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 4, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
return fmt.Errorf("记录 MySQL schema v4 失败: %w", err)
}
}
return CheckMySQLSchema(db)
}
func checkMySQLSchemaV3(db *sql.DB) error {
tables := []string{"shopee_products", "shopee_skus", "pdd_products", "syb_orders", "spec_mappings", "tasks", "clients", "idempotency_keys", "task_claims", "syb_session", "syb_sync_state", "users", "web_sessions", "client_user_assignments", "syb_sync_runs", "admin_initialization_lock"}
if err := checkMySQLSchema(db, tables); err != nil {
return err
}
return checkMySQLV3Shape(db)
}
// migrateMySQLV3 把采购规格身份从蝦皮 SKU 改为顺运宝商品规格原文。
// 每一步都可重放:DDL 已提交但版本尚未记录时,再启动仍会收敛。
func migrateMySQLV3(db *sql.DB) error {
@@ -686,11 +721,44 @@ func CheckMySQLSchema(db *sql.DB) error {
"shopee_products", "shopee_skus", "pdd_products", "syb_orders", "spec_mappings",
"tasks", "clients", "idempotency_keys", "task_claims", "syb_session", "syb_sync_state",
"users", "web_sessions", "client_user_assignments", "syb_sync_runs", "admin_initialization_lock",
"spec_mapping_decisions",
}
if err := checkMySQLSchema(db, mysqlRequiredTables); err != nil {
return err
}
return checkMySQLV3Shape(db)
if err := checkMySQLV3Shape(db); err != nil {
return err
}
return checkMySQLV4Shape(db)
}
func checkMySQLV4Shape(db *sql.DB) error {
for _, c := range []struct {
name string
length int64
nullable bool
coll string
}{
{"shopee_goods_id", 191, false, "utf8mb4_bin"}, {"spec_key", 191, false, "utf8mb4_bin"}, {"pdd_goods_id", 191, false, "utf8mb4_bin"},
{"rules_version", 32, false, "utf8mb4_0900_ai_ci"}, {"suggested_option_key", 191, true, "utf8mb4_bin"}, {"chosen_option_key", 191, false, "utf8mb4_bin"},
} {
if err := checkMySQLVarcharColumn(db, "spec_mapping_decisions", c.name, c.length, c.nullable, c.coll, ""); err != nil {
return err
}
}
var enforced, clause string
if err := db.QueryRow(`SELECT tc.enforced,cc.check_clause FROM information_schema.table_constraints tc JOIN information_schema.check_constraints cc ON cc.constraint_schema=tc.constraint_schema AND cc.constraint_name=tc.constraint_name WHERE tc.constraint_schema=DATABASE() AND tc.table_name='spec_mapping_decisions' AND tc.constraint_name='chk_spec_mapping_decisions_accepted'`).Scan(&enforced, &clause); err != nil {
return fmt.Errorf("审计 accepted CHECK 缺失: %w", err)
}
n := strings.NewReplacer("`", "", " ", "", "(", "", ")", "", "_utf8mb4", "", `\`, "").Replace(strings.ToLower(clause))
if enforced != "YES" || n != "acceptedin0,1" {
return fmt.Errorf("审计 accepted CHECK 不正确")
}
var cols string
if err := db.QueryRow(`SELECT GROUP_CONCAT(column_name ORDER BY seq_in_index) FROM information_schema.statistics WHERE table_schema=DATABASE() AND table_name='spec_mapping_decisions' AND index_name='idx_decisions_mapping'`).Scan(&cols); err != nil || cols != "shopee_goods_id,spec_key,pdd_goods_id" {
return fmt.Errorf("审计映射索引不正确")
}
return nil
}
func checkMySQLSchema(db *sql.DB, tables []string) error {
@@ -183,6 +183,55 @@ func TestMySQLMigrate_V3形状自检失败不记版本(t *testing.T) {
}
}
func TestMySQLMigrate_V3升级V4且断点重跑(t *testing.T) {
db := openMySQLMigrationTestDB(t)
defer db.Close()
cleanMySQLTestSchema(t, db)
defer cleanMySQLTestSchema(t, db)
prepareMySQLV2(t, db)
if err := migrateMySQLV3(db); err != nil {
t.Fatal(err)
}
mustExec(t, db, `INSERT INTO schema_migrations(version,applied_at) VALUES (3,'2026-08-10T00:00:00Z')`)
// 模拟 DDL 已完成、版本未记录。
mustExec(t, db, mysqlSchemaV4Decisions)
if err := MigrateMySQL(db); err != nil {
t.Fatal(err)
}
if err := MigrateMySQL(db); err != nil {
t.Fatalf("重跑失败: %v", err)
}
var versions int
if err := db.QueryRow(`SELECT COUNT(*) FROM schema_migrations WHERE version=4`).Scan(&versions); err != nil || versions != 1 {
t.Fatalf("v4=%d err=%v", versions, err)
}
mustExec(t, db, `INSERT INTO spec_mapping_decisions(shopee_goods_id,spec_key,pdd_goods_id,rules_version,chosen_option_key,accepted,decided_at) VALUES('S','K','P','rules_v1','O',1,'2026-08-10T00:00:00Z')`)
if _, err := db.Exec(`UPDATE spec_mapping_decisions SET accepted=2`); err == nil {
t.Fatal("accepted CHECK 必须拒绝 2")
}
}
func TestMySQLMigrate_V4形状错误不记版本(t *testing.T) {
db := openMySQLMigrationTestDB(t)
defer db.Close()
cleanMySQLTestSchema(t, db)
defer cleanMySQLTestSchema(t, db)
prepareMySQLV2(t, db)
if err := migrateMySQLV3(db); err != nil {
t.Fatal(err)
}
mustExec(t, db, `INSERT INTO schema_migrations(version,applied_at) VALUES (3,'2026-08-10T00:00:00Z')`)
mustExec(t, db, strings.Replace(mysqlSchemaV4Decisions, "CHECK (accepted IN (0,1))", "CHECK (accepted IN (0,1,2))", 1))
if err := MigrateMySQL(db); err == nil {
t.Fatal("错误 CHECK 必须阻止 v4")
}
var count int
db.QueryRow(`SELECT COUNT(*) FROM schema_migrations WHERE version=4`).Scan(&count)
if count != 0 {
t.Fatal("自检失败不得记 v4")
}
}
func openMySQLMigrationTestDB(t *testing.T) *sql.DB {
t.Helper()
if os.Getenv("CMAUTOBUY_MYSQL_TEST") != "1" {
+6 -4
View File
@@ -346,6 +346,7 @@ type SybOrderContext struct {
PddCollectStatus string
PddCollectMsg string
PddSkusJSON string
PddUpdatedAt string
MappingOptionKey string
MappingOptions string
HasActiveTask bool
@@ -356,14 +357,14 @@ func scanSybOrderContext(s rowScanner) (SybOrderContext, error) {
var title, productSpec, specKey, shopeeGoodsID, imageURL sql.NullString
var priceCent sql.NullInt64
var shopeeExists int
var pddGoodsID, pddGoodsURL, collectStatus, collectMsg, skusJSON sql.NullString
var pddGoodsID, pddGoodsURL, collectStatus, collectMsg, skusJSON, pddUpdatedAt sql.NullString
var mappingKey, mappingOptions sql.NullString
var hasActiveTask int
err := s.Scan(
&c.Order.SybID, &c.Order.OrderNo, &title, &productSpec, &specKey, &shopeeGoodsID,
&c.Order.Quantity, &priceCent, &imageURL, &c.Order.SybData,
&c.Order.CreatedAt, &c.Order.UpdatedAt, &shopeeExists,
&pddGoodsID, &pddGoodsURL, &collectStatus, &collectMsg, &skusJSON,
&pddGoodsID, &pddGoodsURL, &collectStatus, &collectMsg, &skusJSON, &pddUpdatedAt,
&mappingKey, &mappingOptions, &hasActiveTask,
)
c.Order.Title = title.String
@@ -378,6 +379,7 @@ func scanSybOrderContext(s rowScanner) (SybOrderContext, error) {
c.PddCollectStatus = collectStatus.String
c.PddCollectMsg = collectMsg.String
c.PddSkusJSON = skusJSON.String
c.PddUpdatedAt = pddUpdatedAt.String
c.MappingOptionKey = mappingKey.String
c.MappingOptions = mappingOptions.String
c.HasActiveTask = hasActiveTask != 0
@@ -393,7 +395,7 @@ func ListSybOrderContexts(q Execer, filter SybOrderFilter, limit, offset int) ([
so.price_twd_cent, so.image_url, so.syb_data, so.created_at, so.updated_at,
CASE WHEN sp.goods_id IS NULL THEN 0 ELSE 1 END,
sp.pdd_goods_id, sp.pdd_goods_url, pp.collect_status, pp.collect_msg,
pp.skus_json, sm.pdd_option_key, sm.pdd_options,
pp.skus_json, pp.updated_at, sm.pdd_option_key, sm.pdd_options,
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
AND t.syb_id = so.syb_id
AND t.status IN ('pending', 'assigned', 'claimed'))` +
@@ -427,7 +429,7 @@ func GetSybOrderContext(q Execer, sybID string) (*SybOrderContext, error) {
so.price_twd_cent, so.image_url, so.syb_data, so.created_at, so.updated_at,
CASE WHEN sp.goods_id IS NULL THEN 0 ELSE 1 END,
sp.pdd_goods_id, sp.pdd_goods_url, pp.collect_status, pp.collect_msg,
pp.skus_json, sm.pdd_option_key, sm.pdd_options,
pp.skus_json, pp.updated_at, sm.pdd_option_key, sm.pdd_options,
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
AND t.syb_id = so.syb_id
AND t.status IN ('pending', 'assigned', 'claimed'))`+