package repository import ( "context" "crypto/tls" "crypto/x509" "database/sql" "fmt" "net" "os" "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, 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) { 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'", } if cfg.TLSMode == config.DatabaseTLSVerifyCA { tlsConfig, err := loadMySQLTLSConfig(cfg.TLSCA) if err != nil { return nil, err } driverConfig.TLS = tlsConfig } return driverConfig, nil } // 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) } roots := x509.NewCertPool() if !roots.AppendCertsFromPEM(pemData) { return nil, fmt.Errorf("解析 MySQL TLS CA 文件 %s 失败: 文件中没有有效 PEM 证书", caPath) } 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 } // 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(¤t); 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 }