Files
cmautobuy/admin/repository/mysql_db.go
T

409 lines
16 KiB
Go

package repository
import (
"context"
"database/sql"
"fmt"
"net"
"strings"
"time"
"github.com/go-sql-driver/mysql"
"cmautobuy/admin/config"
)
const mysqlSchemaVersion = 2
// OpenMySQL 打开生产 MySQL 8 数据库。错误信息绝不包含完整 DSN 或密码。
func OpenMySQL(cfg config.DatabaseConfig) (*sql.DB, error) {
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
// 业务层用 RowsAffected 判断目标行是否存在;重复写入相同值也应算匹配到。
driverConfig.ClientFoundRows = true
driverConfig.Params = map[string]string{
"time_zone": "'+00:00'",
"sql_mode": "'STRICT_TRANS_TABLES,ERROR_FOR_DIVISION_BY_ZERO,NO_ENGINE_SUBSTITUTION'",
}
db, err := sql.Open("mysql", driverConfig.FormatDSN())
if err != nil {
return nil, fmt.Errorf("准备 MySQL 连接失败: %w", err)
}
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 失败,请检查服务、库名和环境变量: %w", err)
}
if err := CheckMySQLServer(db, cfg.Name); err != nil {
db.Close()
return nil, err
}
return db, nil
}
// 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`,
}
// 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)`,
}
// 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)
}
}
if err := checkMySQLSchema(db, requiredTables); err != nil {
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)
}
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)
}
}
if err := CheckMySQLSchema(db); err != nil {
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)
}
}
return CheckMySQLSchema(db)
}
// CheckMySQLSchema 确认所有业务表和关键追加列存在。
func CheckMySQLSchema(db *sql.DB) error {
mysqlRequiredTables := append(append([]string{}, requiredTables...), "admin_initialization_lock")
return checkMySQLSchema(db, mysqlRequiredTables)
}
func checkMySQLSchema(db *sql.DB, tables []string) error {
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
for _, table := range tables {
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
}