feat: 建立 MySQL 8 数据库基础 (#78)

This commit is contained in:
chengma
2026-08-10 01:22:22 +08:00
parent 93bbca88a8
commit 7e74bf210d
12 changed files with 623 additions and 41 deletions
+17 -11
View File
@@ -25,9 +25,13 @@
**Go 版本固定 1.23.0。** 依赖版本受它约束,见下面的表。
- 固定使用 **Go + Gin + Go 标准库 `html/template`**。
- 数据库固定 **SQLite**,驱动固定 `modernc.org/sqlite`。
**不得改用 `mattn/go-sqlite3`** —— 那个要 cgo,Windows 上得装 gcc,
交叉编译和打包 exe 都会变得很麻烦。
- 生产数据库固定 **MySQL 8.4**,驱动固定
`github.com/go-sql-driver/mysql` **v1.9.2**(纯 Go,兼容 Go 1.23)。
- Admin 必须与 MySQL 同机或走私有网络;生产默认连接 `127.0.0.1:3307`,
**不得为了远程访问开放公网 MySQL 端口**。
- `modernc.org/sqlite` 只保留给 SQLite → MySQL 单向迁移工具和历史库回归,
不得用于生产运行时,也不得做 SQLite/MySQL 双写或自动同步。
- 仍然**不得改用 `mattn/go-sqlite3`**,避免引入 cgo 和 gcc。
- Excel 读取固定 `github.com/xuri/excelize/v2`。
### 依赖版本已钉死,不要随手升
@@ -37,12 +41,13 @@
| `github.com/gin-gonic/gin` | **v1.11.0** | v1.12.0 起要求 Go ≥ 1.25.0 |
| `modernc.org/sqlite` | **v1.38.0** | v1.40.0 起要求 Go ≥ 1.24.0,v1.48.0 起要求 ≥ 1.25.0 |
| `github.com/xuri/excelize/v2` | **v2.9.1** | v2.10.0 要求 Go ≥ 1.24.0,v2.11.0 要求 ≥ 1.25.0 |
| `github.com/go-sql-driver/mysql` | **v1.9.2** | v1.9.3 起要求 Go ≥ 1.24.0 |
> excelize 目前**还不在 `go.mod` 里**——导入功能还没写,没有代码 import 它,
> `go mod tidy` 会把它去掉,这是 Go 的正常行为。写导入功能时用
> `go get github.com/xuri/excelize/v2@v2.9.1` 加进来。
这三个版本是**在 Go 1.23.0 下实测能编译通过的最高版本**。
这些版本是**在 Go 1.23.0 下实测能编译通过的固定版本**。
`[必须]` 直接跑 `go get <包名>`(不带版本)会拉到最新版,然后报
`requires go >= 1.25.0`。拉依赖要带版本号,见
@@ -68,7 +73,7 @@
- `handler` 只负责解析请求、调用 service、渲染模板或返回 JSON,**不写业务逻辑,不拼 SQL**。
- `service` 放业务逻辑(导入解析、创建任务、匹配复用),不认识 Gin 的 `*gin.Context`。
- `repository` 封装 SQLite,**只有这一层能写 SQL**。
- `repository` 封装 MySQL 读写,迁移命令中的 SQLite 读取也只能放在 repository/cmd 边界内;**只有这一层能写业务 SQL**。
- `model` 只放数据结构,不导入 Gin 和数据库驱动。
- 给 Client 的接口和给浏览器的页面**分开放**(`handler/api/` 和 `handler/web/`),
两者的错误格式、认证方式都不一样,混在一起迟早出事。
@@ -85,9 +90,9 @@
- `[必须]` 蝦皮规格原文(`spec_raw`)永远保留,解析不出来就留空,不要瞎猜。
- `[必须]` Client 提交结果时,**不管任务是否已取消、是否已重派,一律接受**,
理由见 Client 契约 §6.1。这条最容易被顺手违反。
- `[必须]` `repository/db.go` 里的 `migrations` 只追加,**不得修改已经发布过的条目**。
改了的话,已经建过库的机器 `user_version` 已经越过它,永远不会重跑,
程序会拿着对不上的库静默启动(见 #20)。需要改结构就加新的一条。
- `[必须]` MySQL 的 `schema_migrations` 只追加,已经发布的版本不得改写。
MySQL DDL 会隐式提交:每条 DDL 必须可重放,整版完成并通过 schema 自检后
才记录版本。历史 SQLite migrations 已冻结,只供迁移工具读取旧库。
## 界面规则
@@ -132,9 +137,10 @@ Remove-Item Env:GOTOOLCHAIN
新加依赖时尤其要跑——依赖的 `go` 指令高于 1.23.0 的话,只有这条命令能发现。
- `go run .` 启动后访问 `http://localhost:8080`,**不会连手机、不会下单**,可随时运行。
- 修改数据库时测试首次建库和从上一版本迁移。
`[必须]` **不要只测全新库。** 先列出现实中存在哪些 schema 状态
(不同版本的程序建过的库都算),一个个迁过来验。见 `docs/task/20-*.md`。
- 修改数据库时必须使用独立 MySQL 8.4 测试库验证首次建库、重复迁移和上一版本升级。
测试库名必须以 `_test` 结尾,禁止测试清理连接生产库。
- 修改 SQLite 单向迁移时,同时用各历史 SQLite schema 状态演练;源库只读,
目标库必须为空,逐表计数和关键业务关系必须核对。
- 修改给 Client 的接口时,跑契约测试,确认仍满足 Client 侧 §6.1 的无条件接受。
- 修改 Excel 导入时用 `raw_data/` 下的样本跑一遍,核对导入条数。
该样本含商业数据、**不在仓库里**,需向项目负责人索取;自动化测试用 `testdata/` 下的脱敏小样本。
+60
View File
@@ -23,6 +23,66 @@ import (
// 本项目有意不做心跳,见 docs/admin/04-client-api.md §3。
const OnlineThreshold = 10 * time.Minute
const (
databaseHostEnv = "CMAUTOBUY_DB_HOST"
databasePortEnv = "CMAUTOBUY_DB_PORT"
databaseNameEnv = "CMAUTOBUY_DB_NAME"
databaseUserEnv = "CMAUTOBUY_DB_USER"
databasePasswordEnv = "CMAUTOBUY_DB_PASSWORD"
)
// DatabaseConfig 是生产 MySQL 8 的连接配置。
// 密码只从环境变量读取,不能写进 config.yaml、日志或工单。
type DatabaseConfig struct {
Host string
Port string
Name string
User string
Password string
}
// String 永远隐藏密码,防止排错时用 %v 把凭据写进日志。
func (c DatabaseConfig) String() string {
password := "(空)"
if c.Password != "" {
password = "****"
}
return fmt.Sprintf("DatabaseConfig{Host:%s Port:%s Name:%s User:%s Password:%s}",
c.Host, c.Port, c.Name, c.User, password)
}
// LoadDatabaseFromEnv 读取生产数据库配置。
// Host/Port 使用线上同机部署的安全默认值;库名、账号和密码必须显式提供。
func LoadDatabaseFromEnv() (DatabaseConfig, error) {
cfg := DatabaseConfig{
Host: strings.TrimSpace(os.Getenv(databaseHostEnv)),
Port: strings.TrimSpace(os.Getenv(databasePortEnv)),
Name: strings.TrimSpace(os.Getenv(databaseNameEnv)),
User: strings.TrimSpace(os.Getenv(databaseUserEnv)),
Password: os.Getenv(databasePasswordEnv),
}
if cfg.Host == "" {
cfg.Host = "127.0.0.1"
}
if cfg.Port == "" {
cfg.Port = "3307"
}
var missing []string
if cfg.Name == "" {
missing = append(missing, databaseNameEnv)
}
if cfg.User == "" {
missing = append(missing, databaseUserEnv)
}
if cfg.Password == "" {
missing = append(missing, databasePasswordEnv)
}
if len(missing) > 0 {
return DatabaseConfig{}, fmt.Errorf("缺少 MySQL 配置环境变量:%s", strings.Join(missing, "、"))
}
return cfg, nil
}
// DataDir 返回可写数据目录,不存在就创建。
//
// 打包成 exe 后 = exe 旁边的 data/
+29
View File
@@ -7,6 +7,35 @@ import (
"testing"
)
func TestLoadDatabaseFromEnv_默认只连接本机3307且密码不泄露(t *testing.T) {
t.Setenv(databaseHostEnv, "")
t.Setenv(databasePortEnv, "")
t.Setenv(databaseNameEnv, "cmautobuy_test")
t.Setenv(databaseUserEnv, "cmautobuy_test")
t.Setenv(databasePasswordEnv, "secret-do-not-print")
cfg, err := LoadDatabaseFromEnv()
if err != nil {
t.Fatal(err)
}
if cfg.Host != "127.0.0.1" || cfg.Port != "3307" {
t.Fatalf("默认地址应为线上同机 MySQL 8,实际 %s:%s", cfg.Host, cfg.Port)
}
if got := cfg.String(); strings.Contains(got, cfg.Password) || !strings.Contains(got, "****") {
t.Fatalf("配置字符串泄露密码或没有打码:%s", got)
}
}
func TestLoadDatabaseFromEnv_缺少必填项时只报变量名(t *testing.T) {
t.Setenv(databaseNameEnv, "")
t.Setenv(databaseUserEnv, "")
t.Setenv(databasePasswordEnv, "")
_, err := LoadDatabaseFromEnv()
if err == nil || !strings.Contains(err.Error(), databasePasswordEnv) {
t.Fatalf("应明确列出缺失变量,实际 %v", err)
}
}
// loadFrom 是测试专用的小工具:把一段 YAML 文本写到临时目录里的
// config.yaml,绕开 ConfigPath()(它依赖 os.Executable,测试环境里
// 不可控),直接测 Load 里"读文件 + 解析"这段逻辑。
+4 -1
View File
@@ -5,7 +5,8 @@ go 1.23.0
// 依赖用 go get 添加后会自动写到下面,连同 go.sum 一起提交。
// 本项目固定使用(理由见 admin/AGENTS.md 技术栈一节):
// github.com/gin-gonic/gin Web 框架
// modernc.org/sqlite SQLite 驱动,纯 Go,不需要 cgo
// github.com/go-sql-driver/mysql 生产 MySQL 8 驱动,纯 Go
// modernc.org/sqlite 只读历史 SQLite 迁移驱动,纯 Go
// github.com/xuri/excelize/v2 读蝦皮 Excel 报表
//
// 不得改用 github.com/mattn/go-sqlite3(需要 cgo,Windows 上要装 gcc,
@@ -13,6 +14,7 @@ go 1.23.0
require (
github.com/gin-gonic/gin v1.11.0
github.com/go-sql-driver/mysql v1.9.2
github.com/goccy/go-yaml v1.18.0
github.com/xuri/excelize/v2 v2.9.1
golang.org/x/crypto v0.40.0
@@ -20,6 +22,7 @@ require (
)
require (
filippo.io/edwards25519 v1.1.0 // indirect
github.com/bytedance/sonic v1.14.0 // indirect
github.com/bytedance/sonic/loader v0.3.0 // indirect
github.com/cloudwego/base64x v0.1.6 // indirect
+4
View File
@@ -1,3 +1,5 @@
filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA=
filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4=
github.com/bytedance/sonic v1.14.0 h1:/OfKt8HFw0kh2rj8N0F6C/qPGRESq0BbaNZgcNXXzQQ=
github.com/bytedance/sonic v1.14.0/go.mod h1:WoEbx8WTcFJfzCe0hbmyTGrfjt8PzNEBdxlNUO24NhA=
github.com/bytedance/sonic/loader v0.3.0 h1:dskwH8edlzNMctoruo8FPTJDF3vLtDT0sXZwvZJyqeA=
@@ -23,6 +25,8 @@ github.com/go-playground/universal-translator v0.18.1 h1:Bcnm0ZwsGyWbCzImXv+pAJn
github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY=
github.com/go-playground/validator/v10 v10.27.0 h1:w8+XrWVMhGkxOaaowyKH35gFydVHOvC0/uWoy2Fzwn4=
github.com/go-playground/validator/v10 v10.27.0/go.mod h1:I5QpIEbmr8On7W0TktmJAumgzX4CA1XNl4ZmDuVHKKo=
github.com/go-sql-driver/mysql v1.9.2 h1:4cNKDYQ1I84SXslGddlsrMhc8k4LeDVj6Ad6WRjiHuU=
github.com/go-sql-driver/mysql v1.9.2/go.mod h1:qn46aNg1333BRMNU69Lq93t8du/dwxI64Gl8i5p1WMU=
github.com/goccy/go-json v0.10.2 h1:CrxCmQqYDkv1z7lO7Wbh2HN93uovUHgrECaO5ZrCXAU=
github.com/goccy/go-json v0.10.2/go.mod h1:6MelG93GURQebXPDq3khkgXZkazVtN9CRI+MGFi0w8I=
github.com/goccy/go-yaml v1.18.0 h1:8W7wMFS12Pcas7KU+VVkaiCng+kG8QiFeFwzFb+rwuw=
+375
View File
@@ -0,0 +1,375 @@
package repository
import (
"context"
"database/sql"
"fmt"
"net"
"strings"
"time"
"github.com/go-sql-driver/mysql"
"cmautobuy/admin/config"
)
const mysqlSchemaVersion = 1
// 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
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`,
}
// 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); 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)
}
}
return CheckMySQLSchema(db)
}
// CheckMySQLSchema 确认所有业务表和关键追加列存在。
func CheckMySQLSchema(db *sql.DB) 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 requiredTables {
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
}
@@ -0,0 +1,79 @@
package repository
import (
"database/sql"
"fmt"
"os"
"strings"
"testing"
"cmautobuy/admin/config"
)
// TestMySQLMigrate_真实MySQL8 只在显式提供隔离测试库时运行。
// 库名必须以 _test 结尾,防止测试清理误碰生产库。
func TestMySQLMigrate_真实MySQL8(t *testing.T) {
if os.Getenv("CMAUTOBUY_MYSQL_TEST") != "1" {
t.Skip("未启用真实 MySQL 8 集成测试")
}
cfg, err := config.LoadDatabaseFromEnv()
if err != nil {
t.Fatal(err)
}
if !strings.HasSuffix(cfg.Name, "_test") {
t.Fatalf("拒绝清理非测试数据库 %q:库名必须以 _test 结尾", cfg.Name)
}
db, err := OpenMySQL(cfg)
if err != nil {
t.Fatal(err)
}
defer db.Close()
cleanMySQLTestSchema(t, db)
defer cleanMySQLTestSchema(t, db)
if err := MigrateMySQL(db); err != nil {
t.Fatalf("首次建立 MySQL schema 失败: %v", err)
}
if err := MigrateMySQL(db); err != nil {
t.Fatalf("重复迁移应该无副作用: %v", err)
}
var version int
if err := db.QueryRow(`SELECT MAX(version) FROM schema_migrations`).Scan(&version); err != nil {
t.Fatal(err)
}
if version != mysqlSchemaVersion {
t.Fatalf("schema 版本=%d,期望 %d", version, mysqlSchemaVersion)
}
}
func cleanMySQLTestSchema(t *testing.T, db *sql.DB) {
t.Helper()
rows, err := db.Query(`SELECT table_name FROM information_schema.tables
WHERE table_schema = DATABASE() AND table_type = 'BASE TABLE'`)
if err != nil {
t.Fatal(err)
}
var tables []string
for rows.Next() {
var table string
if err := rows.Scan(&table); err != nil {
rows.Close()
t.Fatal(err)
}
tables = append(tables, table)
}
if err := rows.Close(); err != nil {
t.Fatal(err)
}
if _, err := db.Exec(`SET FOREIGN_KEY_CHECKS = 0`); err != nil {
t.Fatal(err)
}
defer db.Exec(`SET FOREIGN_KEY_CHECKS = 1`)
for _, table := range tables {
// 表名只来自当前 _test 数据库的 information_schema,并对反引号转义。
quoted := "`" + strings.ReplaceAll(table, "`", "``") + "`"
if _, err := db.Exec("DROP TABLE " + quoted); err != nil {
t.Fatal(fmt.Errorf("清理 MySQL 测试表 %s 失败: %w", table, err))
}
}
}