2026-08-10 01:22:22 +08:00
package repository
import (
"context"
2026-08-10 09:59:09 +08:00
"crypto/tls"
"crypto/x509"
2026-08-10 01:22:22 +08:00
"database/sql"
"fmt"
2026-08-10 12:24:23 +08:00
"log"
2026-08-10 01:22:22 +08:00
"net"
2026-08-10 09:59:09 +08:00
"os"
2026-08-10 01:22:22 +08:00
"strings"
"time"
"github.com/go-sql-driver/mysql"
"cmautobuy/admin/config"
2026-08-10 12:24:23 +08:00
"cmautobuy/admin/model"
"cmautobuy/admin/spec"
2026-08-10 01:22:22 +08:00
)
2026-08-10 12:41:17 +08:00
const mysqlSchemaVersion = 4
2026-08-10 01:22:22 +08:00
// OpenMySQL 打开生产 MySQL 8 数据库。错误信息绝不包含完整 DSN 或密码。
func OpenMySQL ( cfg config . DatabaseConfig ) ( * sql . DB , error ) {
2026-08-10 09:59:09 +08:00
driverConfig , err := newMySQLDriverConfig ( cfg )
if err != nil {
return nil , err
}
connector , err := mysql . NewConnector ( driverConfig )
if err != nil {
return nil , fmt . Errorf ( "准备 MySQL 连接失败: %w" , err )
}
db := sql . OpenDB ( connector )
db . SetMaxOpenConns ( 10 )
db . SetMaxIdleConns ( 10 )
db . SetConnMaxLifetime ( 3 * time . Minute )
db . SetConnMaxIdleTime ( time . Minute )
ctx , cancel := context . WithTimeout ( context . Background (), 5 * time . Second )
defer cancel ()
if err := db . PingContext ( ctx ); err != nil {
db . Close ()
return nil , fmt . Errorf ( "连接 MySQL 失败,请检查服务、库名、config.yaml、环境变量和 TLS 配置: %w" , err )
}
if err := CheckMySQLServer ( db , cfg . Name ); err != nil {
db . Close ()
return nil , err
}
return db , nil
}
// newMySQLDriverConfig 只负责把业务配置转换成驱动配置,方便单元测试在不连接
// 数据库的情况下核对 TLS 是否真的开启。
func newMySQLDriverConfig ( cfg config . DatabaseConfig ) ( * mysql . Config , error ) {
2026-08-10 01:22:22 +08:00
driverConfig := mysql . NewConfig ()
driverConfig . User = cfg . User
driverConfig . Passwd = cfg . Password
driverConfig . Net = "tcp"
driverConfig . Addr = net . JoinHostPort ( cfg . Host , cfg . Port )
driverConfig . DBName = cfg . Name
driverConfig . Collation = "utf8mb4_0900_ai_ci"
driverConfig . Loc = time . UTC
driverConfig . Timeout = 5 * time . Second
driverConfig . ReadTimeout = 30 * time . Second
driverConfig . WriteTimeout = 30 * time . Second
driverConfig . RejectReadOnly = true
2026-08-10 02:21:31 +08:00
// 业务层用 RowsAffected 判断目标行是否存在;重复写入相同值也应算匹配到。
driverConfig . ClientFoundRows = true
2026-08-10 01:22:22 +08:00
driverConfig . Params = map [ string ] string {
"time_zone" : "'+00:00'" ,
"sql_mode" : "'STRICT_TRANS_TABLES,ERROR_FOR_DIVISION_BY_ZERO,NO_ENGINE_SUBSTITUTION'" ,
}
2026-08-10 09:59:09 +08:00
if cfg . TLSMode == config . DatabaseTLSVerifyCA {
tlsConfig , err := loadMySQLTLSConfig ( cfg . TLSCA )
if err != nil {
return nil , err
}
driverConfig . TLS = tlsConfig
2026-08-10 01:22:22 +08:00
}
2026-08-10 09:59:09 +08:00
return driverConfig , nil
}
2026-08-10 01:22:22 +08:00
2026-08-10 09:59:09 +08:00
// loadMySQLTLSConfig 加载服务器专属 CA,并验证服务端证书确实由它签发。
// MySQL 自动生成的服务端证书没有 SAN,Go 无法做 IP/域名匹配;这里显式跳过
// 内置主机名检查,但用 VerifyConnection 恢复证书链校验,不能退回明文连接。
func loadMySQLTLSConfig ( caPath string ) ( * tls . Config , error ) {
pemData , err := os . ReadFile ( caPath )
if err != nil {
return nil , fmt . Errorf ( "读取 MySQL TLS CA 文件 %s 失败: %w" , caPath , err )
2026-08-10 01:22:22 +08:00
}
2026-08-10 09:59:09 +08:00
roots := x509 . NewCertPool ()
if ! roots . AppendCertsFromPEM ( pemData ) {
return nil , fmt . Errorf ( "解析 MySQL TLS CA 文件 %s 失败: 文件中没有有效 PEM 证书" , caPath )
2026-08-10 01:22:22 +08:00
}
2026-08-10 09:59:09 +08:00
return & tls . Config {
MinVersion : tls . VersionTLS12 ,
InsecureSkipVerify : true , // 主机名检查由下面的服务器专属 CA 链校验替代。
VerifyConnection : func ( state tls . ConnectionState ) error {
if len ( state . PeerCertificates ) == 0 {
return fmt . Errorf ( "验证 MySQL TLS 证书失败: 服务端没有提供证书" )
}
intermediates := x509 . NewCertPool ()
for _ , certificate := range state . PeerCertificates [ 1 :] {
intermediates . AddCert ( certificate )
}
_ , err := state . PeerCertificates [ 0 ]. Verify ( x509 . VerifyOptions {
Roots : roots ,
Intermediates : intermediates ,
KeyUsages : [] x509 . ExtKeyUsage { x509 . ExtKeyUsageServerAuth },
})
if err != nil {
return fmt . Errorf ( "验证 MySQL TLS 证书链失败: %w" , err )
}
return nil
},
}, nil
2026-08-10 01:22:22 +08:00
}
// CheckMySQLServer 拒绝错误版本、错误库、非 UTC 或非 utf8mb4 的连接。
func CheckMySQLServer ( db * sql . DB , expectedDatabase string ) error {
var version , databaseName , timeZone , charset , collation string
if err := db . QueryRow ( `SELECT VERSION(), DATABASE(), @@session.time_zone,
@@character_set_connection, @@collation_connection` ). Scan (
& version , & databaseName , & timeZone , & charset , & collation ); err != nil {
return fmt . Errorf ( "读取 MySQL 运行参数失败: %w" , err )
}
if ! strings . HasPrefix ( version , "8." ) {
return fmt . Errorf ( "MySQL 版本不兼容:需要 8.x,实际 %s" , version )
}
if databaseName != expectedDatabase {
return fmt . Errorf ( "连接到了错误的 MySQL 数据库:期望 %s,实际 %s" , expectedDatabase , databaseName )
}
if timeZone != "+00:00" {
return fmt . Errorf ( "MySQL 会话时区必须是 +00:00,实际 %s" , timeZone )
}
if charset != "utf8mb4" || ! strings . HasPrefix ( collation , "utf8mb4_" ) {
return fmt . Errorf ( "MySQL 连接字符集必须是 utf8mb4,实际 %s/%s" , charset , collation )
}
return nil
}
// mysqlSchemaV1 是当前 SQLite v8 最终业务结构在 MySQL 8 上的等价定义。
// 每条 CREATE TABLE 都可重复执行;全部完成并通过自检后才记录版本。
var mysqlSchemaV1 = [] string {
`CREATE TABLE IF NOT EXISTS shopee_products (
goods_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
title TEXT NOT NULL,
shopee_status VARCHAR(191),
main_sku_code VARCHAR(191),
pdd_goods_url TEXT,
pdd_goods_id VARCHAR(191) COLLATE utf8mb4_bin,
created_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL,
KEY idx_shopee_products_pdd (pdd_goods_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS shopee_skus (
sku_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
goods_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
spec_raw TEXT NOT NULL,
color TEXT,
size TEXT,
advice TEXT,
parse_ok TINYINT NOT NULL DEFAULT 0,
sku_code VARCHAR(191),
is_manual TINYINT NOT NULL DEFAULT 0,
created_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL,
KEY idx_shopee_skus_goods (goods_id),
KEY idx_shopee_skus_parse (parse_ok),
CONSTRAINT fk_shopee_skus_product FOREIGN KEY (goods_id)
REFERENCES shopee_products(goods_id) ON DELETE CASCADE,
CHECK (parse_ok IN (0, 1)),
CHECK (is_manual IN (0, 1))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS pdd_products (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
goods_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL UNIQUE,
url TEXT NOT NULL,
title TEXT,
shop_name TEXT,
skus_json LONGTEXT,
collect_status VARCHAR(20) NOT NULL DEFAULT 'pending',
collect_msg TEXT,
artifact_ref TEXT,
collected_at VARCHAR(35),
deleted_at VARCHAR(35),
created_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL,
KEY idx_pdd_products_status (collect_status),
CHECK (collect_status IN ('pending', 'collecting', 'collected', 'failed'))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS syb_orders (
syb_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
order_no VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
title TEXT,
shopee_goods_id VARCHAR(191) COLLATE utf8mb4_bin,
shopee_sku_id VARCHAR(191) COLLATE utf8mb4_bin,
product_spec TEXT,
quantity BIGINT NOT NULL,
price_twd_cent BIGINT,
image_url TEXT,
syb_data LONGTEXT NOT NULL DEFAULT ('{}'),
created_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL,
KEY idx_syb_orders_order (order_no),
KEY idx_syb_orders_goods (shopee_goods_id),
KEY idx_syb_orders_list (updated_at DESC, syb_id DESC),
CHECK (quantity > 0),
CHECK (price_twd_cent IS NULL OR price_twd_cent >= 0)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS sku_mappings (
shopee_sku_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
pdd_goods_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
pdd_option_key VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
pdd_options LONGTEXT NOT NULL,
goods_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
mapped_at VARCHAR(35) NOT NULL,
mapped_by VARCHAR(191),
PRIMARY KEY (shopee_sku_id, pdd_goods_id),
KEY idx_sku_mappings_goods (goods_id),
KEY idx_sku_mappings_pdd (pdd_goods_id),
CONSTRAINT fk_sku_mappings_sku FOREIGN KEY (shopee_sku_id)
REFERENCES shopee_skus(sku_id) ON DELETE CASCADE
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS tasks (
task_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
task_type VARCHAR(20) NOT NULL,
status VARCHAR(20) NOT NULL DEFAULT 'pending',
version BIGINT NOT NULL DEFAULT 1,
priority BIGINT NOT NULL DEFAULT 0,
assigned_client VARCHAR(191) COLLATE utf8mb4_bin,
claimed_at VARCHAR(35),
syb_id VARCHAR(191) COLLATE utf8mb4_bin,
order_no VARCHAR(191) COLLATE utf8mb4_bin,
goods_id VARCHAR(191) COLLATE utf8mb4_bin,
shopee_sku_id VARCHAR(191) COLLATE utf8mb4_bin,
pdd_goods_url TEXT NOT NULL,
pdd_goods_id VARCHAR(191) COLLATE utf8mb4_bin,
pdd_options LONGTEXT,
quantity BIGINT,
max_price_cent BIGINT,
result_data LONGTEXT,
error_code VARCHAR(191),
error_message TEXT,
finished_at VARCHAR(35),
created_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL,
KEY idx_tasks_claim (assigned_client, status, priority DESC, created_at),
KEY idx_tasks_list (updated_at DESC, task_id DESC),
KEY idx_tasks_order (order_no),
CHECK (task_type IN ('collect', 'purchase')),
CHECK (status IN ('pending', 'assigned', 'claimed', 'succeeded', 'manual_review', 'failed', 'cancelled')),
CHECK (version > 0),
CHECK (quantity IS NULL OR quantity > 0),
CHECK (max_price_cent IS NULL OR max_price_cent > 0)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS clients (
client_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
name TEXT,
device_address TEXT,
platform VARCHAR(191),
pdd_package VARCHAR(191),
capabilities LONGTEXT,
last_seen_at VARCHAR(35) NOT NULL,
created_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS idempotency_keys (
` + "`key`" + ` VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
request_hash VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
response_body LONGTEXT NOT NULL,
created_at VARCHAR(35) NOT NULL
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS task_claims (
task_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
client_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
claimed_at VARCHAR(35) NOT NULL,
PRIMARY KEY (task_id, client_id),
KEY idx_task_claims_client (client_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS syb_session (
username VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
cookies LONGTEXT NOT NULL,
expires_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS syb_sync_state (
id TINYINT PRIMARY KEY,
last_synced_at VARCHAR(35),
updated_at VARCHAR(35) NOT NULL,
CHECK (id = 1)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS users (
user_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
username VARCHAR(191) COLLATE utf8mb4_0900_ai_ci NOT NULL UNIQUE,
password_hash VARCHAR(255) COLLATE utf8mb4_bin NOT NULL,
role VARCHAR(20) NOT NULL,
status VARCHAR(20) NOT NULL,
last_login_at VARCHAR(35),
password_changed_at VARCHAR(35) NOT NULL,
created_at VARCHAR(35) NOT NULL,
updated_at VARCHAR(35) NOT NULL,
CHECK (role IN ('admin', 'purchaser')),
CHECK (status IN ('active', 'disabled'))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS web_sessions (
session_hash VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
user_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
expires_at VARCHAR(35) NOT NULL,
created_at VARCHAR(35) NOT NULL,
last_seen_at VARCHAR(35) NOT NULL,
KEY idx_web_sessions_user (user_id),
KEY idx_web_sessions_expiry (expires_at),
CONSTRAINT fk_web_sessions_user FOREIGN KEY (user_id) REFERENCES users(user_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS client_user_assignments (
assignment_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
client_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
user_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
started_at VARCHAR(35) NOT NULL,
ended_at VARCHAR(35),
assigned_by_user_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
ended_by_user_id VARCHAR(191) COLLATE utf8mb4_bin,
end_reason VARCHAR(20),
current_client_id VARCHAR(191) COLLATE utf8mb4_bin
GENERATED ALWAYS AS (IF(ended_at IS NULL, client_id, NULL)) STORED,
UNIQUE KEY idx_client_assignment_current (current_client_id),
KEY idx_client_assignment_user (user_id, ended_at, client_id),
KEY idx_client_assignment_history (client_id, started_at DESC),
CONSTRAINT fk_client_assignment_user FOREIGN KEY (user_id) REFERENCES users(user_id),
CONSTRAINT fk_client_assignment_assigned_by FOREIGN KEY (assigned_by_user_id) REFERENCES users(user_id),
CONSTRAINT fk_client_assignment_ended_by FOREIGN KEY (ended_by_user_id) REFERENCES users(user_id),
CHECK (end_reason IS NULL OR end_reason IN ('unbind', 'transfer')),
CHECK ((ended_at IS NULL AND ended_by_user_id IS NULL AND end_reason IS NULL)
OR (ended_at IS NOT NULL AND ended_by_user_id IS NOT NULL AND end_reason IS NOT NULL))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`CREATE TABLE IF NOT EXISTS syb_sync_runs (
run_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
user_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
date_from VARCHAR(10) NOT NULL,
date_to VARCHAR(10) NOT NULL,
status VARCHAR(20) NOT NULL,
stock_count BIGINT NOT NULL DEFAULT 0,
detail_count BIGINT NOT NULL DEFAULT 0,
created_count BIGINT NOT NULL DEFAULT 0,
updated_count BIGINT NOT NULL DEFAULT 0,
skipped_count BIGINT NOT NULL DEFAULT 0,
error_message TEXT,
cursor_advanced TINYINT NOT NULL DEFAULT 0,
started_at VARCHAR(35) NOT NULL,
finished_at VARCHAR(35),
KEY idx_syb_sync_runs_started (started_at DESC, run_id DESC),
CONSTRAINT fk_syb_sync_runs_user FOREIGN KEY (user_id) REFERENCES users(user_id),
CHECK (status IN ('running', 'succeeded', 'failed', 'interrupted')),
CHECK (cursor_advanced IN (0, 1)),
CHECK ((status = 'running' AND finished_at IS NULL)
OR (status <> 'running' AND finished_at IS NOT NULL))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
}
2026-08-10 02:21:31 +08:00
// mysqlSchemaV2 增加首位管理员初始化的并发哨兵。MySQL DDL 会隐式提交,
// 因此每条语句都必须可重放,并在全部成功后才记录版本。
var mysqlSchemaV2 = [] string {
`CREATE TABLE IF NOT EXISTS admin_initialization_lock (
id TINYINT PRIMARY KEY,
CHECK (id = 1)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ,
`INSERT INTO admin_initialization_lock (id) VALUES (1)
ON DUPLICATE KEY UPDATE id = VALUES(id)` ,
}
2026-08-10 12:24:23 +08:00
const mysqlSchemaV3SpecMappings = `CREATE TABLE IF NOT EXISTS spec_mappings (
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,
pdd_option_key VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
pdd_options LONGTEXT NOT NULL,
spec_raw TEXT NOT NULL,
mapped_at VARCHAR(35) NOT NULL,
mapped_by VARCHAR(191),
PRIMARY KEY (shopee_goods_id, spec_key, pdd_goods_id),
KEY idx_spec_mappings_goods (shopee_goods_id),
KEY idx_spec_mappings_pdd (pdd_goods_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`
2026-08-10 12:41:17 +08:00
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`
2026-08-10 01:22:22 +08:00
// MigrateMySQL 建立或升级 MySQL schema。生产迁移只能在这里追加新版本。
func MigrateMySQL ( db * sql . DB ) error {
if _ , err := db . Exec ( `CREATE TABLE IF NOT EXISTS schema_migrations (
version INT PRIMARY KEY,
applied_at VARCHAR(35) NOT NULL
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci` ); err != nil {
return fmt . Errorf ( "建立 MySQL 迁移版本表失败: %w" , err )
}
var current int
if err := db . QueryRow ( `SELECT COALESCE(MAX(version), 0) FROM schema_migrations` ). Scan ( & current ); err != nil {
return fmt . Errorf ( "读取 MySQL schema 版本失败: %w" , err )
}
if current > mysqlSchemaVersion {
return fmt . Errorf ( "MySQL schema 版本 %d 高于程序支持的 %d,请升级程序" , current , mysqlSchemaVersion )
}
if current < 1 {
for i , statement := range mysqlSchemaV1 {
if _ , err := db . Exec ( statement ); err != nil {
return fmt . Errorf ( "执行 MySQL schema v1 第 %d 条失败: %w" , i + 1 , err )
}
}
2026-08-10 02:21:31 +08:00
if err := checkMySQLSchema ( db , requiredTables ); err != nil {
2026-08-10 01:22:22 +08:00
return fmt . Errorf ( "MySQL schema v1 自检失败,未记录版本: %w" , err )
}
if _ , err := db . Exec ( `INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)` ,
1 , time . Now (). UTC (). Format ( time . RFC3339Nano )); err != nil {
return fmt . Errorf ( "记录 MySQL schema v1 失败: %w" , err )
}
2026-08-10 02:21:31 +08:00
current = 1
}
if current < 2 {
for i , statement := range mysqlSchemaV2 {
if _ , err := db . Exec ( statement ); err != nil {
return fmt . Errorf ( "执行 MySQL schema v2 第 %d 条失败: %w" , i + 1 , err )
}
}
2026-08-10 12:24:23 +08:00
mysqlRequiredTablesV2 := append ( append ([] string {}, requiredTables ... ), "admin_initialization_lock" )
if err := checkMySQLSchema ( db , mysqlRequiredTablesV2 ); err != nil {
2026-08-10 02:21:31 +08:00
return fmt . Errorf ( "MySQL schema v2 自检失败,未记录版本: %w" , err )
}
if _ , err := db . Exec ( `INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)` ,
2 , time . Now (). UTC (). Format ( time . RFC3339Nano )); err != nil {
return fmt . Errorf ( "记录 MySQL schema v2 失败: %w" , err )
}
2026-08-10 12:24:23 +08:00
current = 2
}
if current < 3 {
if err := migrateMySQLV3 ( db ); err != nil {
return fmt . Errorf ( "执行 MySQL schema v3 失败: %w" , err )
}
2026-08-10 12:41:17 +08:00
if err := checkMySQLSchemaV3 ( db ); err != nil {
2026-08-10 12:24:23 +08:00
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 )
}
2026-08-10 12:41:17 +08:00
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 )
}
2026-08-10 01:22:22 +08:00
}
return CheckMySQLSchema ( db )
}
2026-08-10 12:41:17 +08:00
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 )
}
2026-08-10 12:24:23 +08:00
// migrateMySQLV3 把采购规格身份从蝦皮 SKU 改为顺运宝商品规格原文。
// 每一步都可重放:DDL 已提交但版本尚未记录时,再启动仍会收敛。
func migrateMySQLV3 ( db * sql . DB ) error {
if err := ensureMySQLV3Columns ( db ); err != nil {
return err
}
if _ , err := db . Exec ( mysqlSchemaV3SpecMappings ); err != nil {
return fmt . Errorf ( "建立 spec_mappings 失败: %w" , err )
}
if err := backfillSybSpecKeys ( db ); err != nil {
return err
}
if err := ensureSybShopeeSkeletons ( db ); err != nil {
return err
}
return convertLegacySKUMappings ( db )
}
func ensureMySQLV3Columns ( db * sql . DB ) error {
exists , err := mysqlColumnExists ( db , "syb_orders" , "spec_key" )
if err != nil {
return err
}
if ! exists {
if _ , err := db . Exec ( `ALTER TABLE syb_orders ADD COLUMN spec_key VARCHAR(191) COLLATE utf8mb4_bin NULL AFTER product_spec` ); err != nil {
return fmt . Errorf ( "增加 syb_orders.spec_key 失败: %w" , err )
}
} else if _ , err := db . Exec ( `ALTER TABLE syb_orders MODIFY COLUMN spec_key VARCHAR(191) COLLATE utf8mb4_bin NULL` ); err != nil {
return fmt . Errorf ( "校正 syb_orders.spec_key 结构失败: %w" , err )
}
exists , err = mysqlColumnExists ( db , "shopee_products" , "source" )
if err != nil {
return err
}
if ! exists {
if _ , err := db . Exec ( `ALTER TABLE shopee_products ADD COLUMN source VARCHAR(16) COLLATE utf8mb4_bin NOT NULL DEFAULT 'report' AFTER main_sku_code` ); err != nil {
return fmt . Errorf ( "增加 shopee_products.source 失败: %w" , err )
}
} else {
// 兼容“列已加但还是可空”的 DDL 中断点。先回填旧行,
// 再收紧 NOT NULL,否则带 NULL 的存量数据会让 ALTER 失败。
if _ , err := db . Exec ( `UPDATE shopee_products SET source='report' WHERE source IS NULL OR source=''` ); err != nil {
return fmt . Errorf ( "回填 shopee_products.source 失败: %w" , err )
}
if _ , err := db . Exec ( `ALTER TABLE shopee_products MODIFY COLUMN source VARCHAR(16) COLLATE utf8mb4_bin NOT NULL DEFAULT 'report'` ); err != nil {
return fmt . Errorf ( "校正 shopee_products.source 结构失败: %w" , err )
}
}
constraintExists , err := mysqlConstraintExists ( db , "shopee_products" , "chk_shopee_products_source" )
if err != nil {
return err
}
if ! constraintExists {
if _ , err := db . Exec ( `ALTER TABLE shopee_products ADD CONSTRAINT chk_shopee_products_source CHECK (source IN ('report','syb'))` ); err != nil {
return fmt . Errorf ( "增加 shopee_products.source CHECK 失败: %w" , err )
}
}
return nil
}
func mysqlColumnExists ( db * sql . DB , table , column string ) ( bool , error ) {
var count int
if err := db . QueryRow ( `SELECT COUNT(*) FROM information_schema.columns
WHERE table_schema=DATABASE() AND table_name=? AND column_name=?` , table , column ). Scan ( & count ); err != nil {
return false , fmt . Errorf ( "检查 MySQL 数据表 %s.%s 失败: %w" , table , column , err )
}
return count == 1 , nil
}
func mysqlConstraintExists ( db * sql . DB , table , constraint string ) ( bool , error ) {
var count int
if err := db . QueryRow ( `SELECT COUNT(*) FROM information_schema.table_constraints
WHERE constraint_schema=DATABASE() AND table_name=? AND constraint_name=?` , table , constraint ). Scan ( & count ); err != nil {
return false , fmt . Errorf ( "检查 MySQL 约束 %s.%s 失败: %w" , table , constraint , err )
}
return count == 1 , nil
}
func backfillSybSpecKeys ( db * sql . DB ) error {
const batchSize = 500
cursor := ""
for {
rows , err := db . Query ( `SELECT syb_id, product_spec FROM syb_orders
WHERE syb_id > ? AND spec_key IS NULL AND product_spec IS NOT NULL
AND TRIM(product_spec) <> '' ORDER BY syb_id LIMIT ?` , cursor , batchSize )
if err != nil {
return fmt . Errorf ( "读取待回填顺运宝规格失败: %w" , err )
}
type item struct { id , raw , key string }
items := make ([] item , 0 , batchSize )
for rows . Next () {
var value item
if err := rows . Scan ( & value . id , & value . raw ); err != nil {
rows . Close ()
return fmt . Errorf ( "读取待回填顺运宝规格失败: %w" , err )
}
value . key , err = spec . SpecKey ( value . raw )
if err != nil {
rows . Close ()
return fmt . Errorf ( "顺运宝明细 %s 的规格不能生成身份键: %w" , value . id , err )
}
items = append ( items , value )
}
if err := rows . Close (); err != nil {
return fmt . Errorf ( "关闭顺运宝规格回填结果失败: %w" , err )
}
if len ( items ) == 0 {
return nil
}
tx , err := db . Begin ()
if err != nil {
return fmt . Errorf ( "开始顺运宝规格回填事务失败: %w" , err )
}
for _ , value := range items {
if _ , err := tx . Exec ( `UPDATE syb_orders SET spec_key=? WHERE syb_id=? AND spec_key IS NULL` , value . key , value . id ); err != nil {
tx . Rollback ()
return fmt . Errorf ( "回填顺运宝明细 %s 的规格键失败: %w" , value . id , err )
}
}
if err := tx . Commit (); err != nil {
return fmt . Errorf ( "提交顺运宝规格回填失败: %w" , err )
}
cursor = items [ len ( items ) - 1 ]. id
}
}
func ensureSybShopeeSkeletons ( db * sql . DB ) error {
now := model . NowISO ()
_ , err := db . Exec ( `INSERT INTO shopee_products
(goods_id, title, source, created_at, updated_at)
SELECT shopee_goods_id, MAX(COALESCE(title,'')), 'syb', ?, ?
FROM syb_orders
WHERE shopee_goods_id IS NOT NULL AND shopee_goods_id <> ''
GROUP BY shopee_goods_id
ON DUPLICATE KEY UPDATE goods_id=VALUES(goods_id)` , now , now )
if err != nil {
return fmt . Errorf ( "按顺运宝数据补建蝦皮商品骨架失败: %w" , err )
}
return nil
}
func convertLegacySKUMappings ( db * sql . DB ) error {
exists , err := mysqlTableExists ( db , "sku_mappings" )
if err != nil || ! exists {
return err
}
if _ , err := db . Exec ( `CREATE TABLE IF NOT EXISTS sku_mappings_v3_backup LIKE sku_mappings` ); err != nil {
return fmt . Errorf ( "建立旧规格映射备份表失败: %w" , err )
}
if _ , err := db . Exec ( `INSERT IGNORE INTO sku_mappings_v3_backup SELECT * FROM sku_mappings` ); err != nil {
return fmt . Errorf ( "备份旧规格映射失败: %w" , err )
}
rows , err := db . Query ( `SELECT m.shopee_sku_id, m.pdd_goods_id, m.pdd_option_key,
m.pdd_options, m.mapped_at, m.mapped_by, sk.goods_id, sk.spec_raw
FROM sku_mappings m LEFT JOIN shopee_skus sk ON sk.sku_id=m.shopee_sku_id` )
if err != nil {
return fmt . Errorf ( "读取旧规格映射失败: %w" , err )
}
type legacy struct {
skuID , pddID , optionKey , options , mappedAt string
mappedBy , goodsID , raw sql . NullString
}
var records [] legacy
for rows . Next () {
var value legacy
if err := rows . Scan ( & value . skuID , & value . pddID , & value . optionKey , & value . options ,
& value . mappedAt , & value . mappedBy , & value . goodsID , & value . raw ); err != nil {
rows . Close ()
return fmt . Errorf ( "读取旧规格映射失败: %w" , err )
}
records = append ( records , value )
}
if err := rows . Close (); err != nil {
return fmt . Errorf ( "关闭旧规格映射结果失败: %w" , err )
}
matched , noCurrentMatch , unable := 0 , 0 , 0
for _ , value := range records {
if ! value . goodsID . Valid || ! value . raw . Valid {
unable ++
log . Printf ( "mysql_v3_mapping_unconvertible shopee_sku_id=%s pdd_goods_id=%s" , value . skuID , value . pddID )
continue
}
key , err := spec . SpecKey ( value . raw . String )
if err != nil {
unable ++
log . Printf ( "mysql_v3_mapping_unconvertible shopee_sku_id=%s pdd_goods_id=%s reason=invalid_spec_key" , value . skuID , value . pddID )
continue
}
if err := UpsertSpecMapping ( db , model . SpecMapping {
ShopeeGoodsID : value . goodsID . String , SpecKey : key , SpecRaw : value . raw . String ,
PddGoodsID : value . pddID , PddOptionKey : value . optionKey , PddOptions : value . options ,
MappedAt : value . mappedAt , MappedBy : value . mappedBy . String ,
}); err != nil {
return fmt . Errorf ( "转换旧规格映射 %s 失败: %w" , value . skuID , err )
}
var hit int
if err := db . QueryRow ( `SELECT EXISTS(SELECT 1 FROM syb_orders WHERE shopee_goods_id=? AND spec_key=?)` ,
value . goodsID . String , key ). Scan ( & hit ); err != nil {
return fmt . Errorf ( "验证旧规格映射 %s 命中失败: %w" , value . skuID , err )
}
if hit == 1 {
matched ++
} else {
noCurrentMatch ++
log . Printf ( "mysql_v3_mapping_no_current_hit shopee_sku_id=%s pdd_goods_id=%s" , value . skuID , value . pddID )
}
}
if matched + noCurrentMatch + unable != len ( records ) {
return fmt . Errorf ( "旧规格映射迁移计数不守恒: 总数=%d 命中=%d 无命中=%d 无法转换=%d" ,
len ( records ), matched , noCurrentMatch , unable )
}
log . Printf ( "mysql_v3_mapping_summary total=%d matched=%d no_current_hit=%d unconvertible=%d backup_table=sku_mappings_v3_backup" ,
len ( records ), matched , noCurrentMatch , unable )
if _ , err := db . Exec ( `DROP TABLE sku_mappings` ); err != nil {
return fmt . Errorf ( "删除已备份的旧规格映射表失败: %w" , err )
}
return nil
}
func mysqlTableExists ( db * sql . DB , table string ) ( bool , error ) {
var count int
if err := db . QueryRow ( `SELECT COUNT(*) FROM information_schema.tables
WHERE table_schema=DATABASE() AND table_name=? AND table_type='BASE TABLE'` , table ). Scan ( & count ); err != nil {
return false , fmt . Errorf ( "检查 MySQL 数据表 %s 失败: %w" , table , err )
}
return count == 1 , nil
}
// CheckMySQLSchema 确认所有业务表、关键追加列和采购身份结构完整。
2026-08-10 01:22:22 +08:00
func CheckMySQLSchema ( db * sql . DB ) error {
2026-08-10 12:24:23 +08:00
mysqlRequiredTables := [] 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" ,
2026-08-10 12:41:17 +08:00
"spec_mapping_decisions" ,
2026-08-10 12:24:23 +08:00
}
if err := checkMySQLSchema ( db , mysqlRequiredTables ); err != nil {
return err
}
2026-08-10 12:41:17 +08:00
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
2026-08-10 02:21:31 +08:00
}
func checkMySQLSchema ( db * sql . DB , tables [] string ) error {
2026-08-10 01:22:22 +08:00
rows , err := db . Query ( `SELECT table_name FROM information_schema.tables
WHERE table_schema = DATABASE() AND table_type = 'BASE TABLE'` )
if err != nil {
return fmt . Errorf ( "读取 MySQL 数据表清单失败: %w" , err )
}
defer rows . Close ()
existing := map [ string ] bool {}
for rows . Next () {
var name string
if err := rows . Scan ( & name ); err != nil {
return fmt . Errorf ( "读取 MySQL 数据表清单失败: %w" , err )
}
existing [ name ] = true
}
if err := rows . Err (); err != nil {
return err
}
var missing [] string
2026-08-10 02:21:31 +08:00
for _ , table := range tables {
2026-08-10 01:22:22 +08:00
if ! existing [ table ] {
missing = append ( missing , table )
}
}
if len ( missing ) > 0 {
return fmt . Errorf ( "MySQL schema 不完整,缺少表 %s" , strings . Join ( missing , "、" ))
}
for table , columns := range requiredColumns {
for _ , column := range columns {
var count int
if err := db . QueryRow ( `SELECT COUNT(*) FROM information_schema.columns
WHERE table_schema = DATABASE() AND table_name = ? AND column_name = ?` ,
table , column ). Scan ( & count ); err != nil {
return fmt . Errorf ( "检查 MySQL 数据表 %s 的列失败: %w" , table , err )
}
if count != 1 {
return fmt . Errorf ( "MySQL 数据表 %s 缺少列 %s" , table , column )
}
}
}
return nil
}
2026-08-10 12:24:23 +08:00
func checkMySQLV3Shape ( db * sql . DB ) error {
rows , err := db . Query ( `SELECT column_name FROM information_schema.key_column_usage
WHERE table_schema=DATABASE() AND table_name='spec_mappings'
AND constraint_name='PRIMARY' ORDER BY ordinal_position` )
if err != nil {
return fmt . Errorf ( "检查 spec_mappings 主键失败: %w" , err )
}
var primary [] string
for rows . Next () {
var column string
if err := rows . Scan ( & column ); err != nil {
rows . Close ()
return fmt . Errorf ( "检查 spec_mappings 主键失败: %w" , err )
}
primary = append ( primary , column )
}
rows . Close ()
if strings . Join ( primary , "," ) != "shopee_goods_id,spec_key,pdd_goods_id" {
return fmt . Errorf ( "spec_mappings 主键不正确,实际为 (%s)" , strings . Join ( primary , "," ))
}
for _ , column := range [] string { "shopee_goods_id" , "spec_key" , "pdd_goods_id" } {
if err := checkMySQLVarcharColumn ( db , "spec_mappings" , column , 191 , false , "utf8mb4_bin" , "" ); err != nil {
return err
}
}
if err := checkMySQLVarcharColumn ( db , "syb_orders" , "spec_key" , 191 , true , "utf8mb4_bin" , "" ); err != nil {
return err
}
if err := checkMySQLVarcharColumn ( db , "shopee_products" , "source" , 16 , false , "utf8mb4_bin" , "report" ); err != nil {
return err
}
var enforced , checkClause 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='shopee_products'
AND tc.constraint_name='chk_shopee_products_source' AND tc.constraint_type='CHECK'` ). Scan ( & enforced , & checkClause ); err != nil {
return fmt . Errorf ( "shopee_products.source CHECK 缺失或不可读: %w" , err )
}
if enforced != "YES" {
return fmt . Errorf ( "shopee_products.source CHECK 未启用" )
}
normalizedClause := strings . NewReplacer ( "`" , "" , " " , "" , "(" , "" , ")" , "" , "_utf8mb4" , "" , `\` , "" ). Replace ( strings . ToLower ( checkClause ))
if ! strings . Contains ( normalizedClause , "sourcein'report','syb'" ) {
return fmt . Errorf ( "shopee_products.source CHECK 表达式不正确" )
}
return nil
}
func checkMySQLVarcharColumn ( db * sql . DB , table , column string , length int64 , nullable bool , collation , defaultValue string ) error {
var dataType , isNullable string
var actualLength sql . NullInt64
var actualCollation , actualDefault sql . NullString
if err := db . QueryRow ( `SELECT data_type, character_maximum_length, is_nullable, collation_name, column_default
FROM information_schema.columns WHERE table_schema=DATABASE() AND table_name=? AND column_name=?` ,
table , column ). Scan ( & dataType , & actualLength , & isNullable , & actualCollation , & actualDefault ); err != nil {
return fmt . Errorf ( "检查 MySQL 列 %s.%s 结构失败: %w" , table , column , err )
}
wantNullable := "NO"
if nullable {
wantNullable = "YES"
}
if dataType != "varchar" || ! actualLength . Valid || actualLength . Int64 != length ||
isNullable != wantNullable || actualCollation . String != collation {
return fmt . Errorf ( "MySQL 列 %s.%s 结构不正确" , table , column )
}
if defaultValue != "" && (! actualDefault . Valid || actualDefault . String != defaultValue ) {
return fmt . Errorf ( "MySQL 列 %s.%s 默认值不正确" , table , column )
}
return nil
}