Files
cmautobuy/admin/repository/db.go
T

1000 lines
42 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package repository 封装 SQLite 读写。
//
// 改动本文件前必读 admin/AGENTS.md。三条硬规则:
// - 只有本包能写 SQL,handler 和 service 都不许拼 SQL;
// - SQL 一律参数化查询(用 ? 占位),禁止字符串拼接;
// - 迁移只能往前加,不许在启动时删库重建——data/ 在升级时是保留的,
// 里面有人工填了几个月的 PDD 链接和 SKU 映射。
//
// 表结构的权威定义在 docs/admin/03-data-model.md,改表要先改文档。
package repository
import (
"context"
"database/sql"
"fmt"
"log"
"path/filepath"
"strings"
// 纯 Go 的 SQLite 驱动,注册的驱动名是 "sqlite"(不是 "sqlite3")。
// 不得换成 github.com/mattn/go-sqlite3,那个需要 cgo,
// Windows 上要装 gcc,打包 exe 会变麻烦。理由见 admin/AGENTS.md。
_ "modernc.org/sqlite"
)
// Execer 让同一个 repository 函数既能直接用 *sql.DB,
// 也能在事务里用 *sql.Tx。
//
// 需要"多张表要么一起改、要么都不改"时,service 开一个事务,
// 把 *sql.Tx 传进来即可,不用为事务再写一套函数。
type Execer interface {
Exec(query string, args ...any) (sql.Result, error)
Query(query string, args ...any) (*sql.Rows, error)
QueryRow(query string, args ...any) *sql.Row
}
// Open 打开 data/admin.db。
//
// **PRAGMA 必须写在 DSN 里,不能用 db.Exec("PRAGMA ...") 设置。**
//
// 原因是 Go 的 database/sql 是一个**连接池**:db.Exec 只作用于当时
// 拿到的那一条连接,池子后来新开的连接完全没执行过那些 PRAGMA。
// 并发写的时候,没有 busy_timeout 的那些连接会直接报
// "database is locked (SQLITE_BUSY)",而不是等锁释放。
//
// 写进 DSN 后,驱动会对**每一条新连接**都应用一遍,这才是对的。
func Open(dataDir string) (*sql.DB, error) {
path := filepath.Join(dataDir, "admin.db")
// 三条设置的含义见 docs/admin/03-data-model.md §2:
// busy_timeout 拿不到锁时最多等 5 秒,而不是立刻报错
// journal_mode WAL 模式,读和写可以同时进行
// foreign_keys 打开外键约束(SQLite 默认是关的)
//
// _txlock=immediate 是另一个**必须加**的参数,原因见下。
dsn := "file:" + path +
"?_pragma=busy_timeout(5000)" +
"&_pragma=journal_mode(WAL)" +
"&_pragma=foreign_keys(1)" +
"&_txlock=immediate"
db, err := sql.Open("sqlite", dsn)
if err != nil {
return nil, fmt.Errorf("打开数据库 %s 失败: %w", path, err)
}
// SQLite 同一时刻只允许一个写事务。连接数放太开,
// 大量连接会互相抢锁、把 busy_timeout 耗光。
// 本项目是单机内部工具,并发量很小,限制在个位数足够。
//
// 注意**不要设成 1**:那样在一个事务里再调用需要连接的代码会死锁。
db.SetMaxOpenConns(4)
db.SetMaxIdleConns(4)
// 关于 _txlock=immediate:
//
// Go 的 db.Begin() 默认发的是 BEGIN DEFERRED——事务开始时**不拿写锁**,
// 等到第一次写才去拿。于是多个事务可以同时开始、各自先读,
// 然后同时想升级成写,互相卡死,直接报 SQLITE_BUSY。
// 这种情况 busy_timeout **救不了**:等下去也不可能有结果,
// 只能让某个事务整个重来。
//
// 加上 _txlock=immediate 后,事务一开始就拿写锁,
// 拿不到就按 busy_timeout 排队等——这才是我们要的行为。
//
// 实测(6 个并发事务,每个先读后写):
// 默认 deferred 失败 5/6
// _txlock=immediate 失败 0/6
// sql.Open 是懒加载的,这里主动连一次,好让配置错误立刻暴露
if err := db.Ping(); err != nil {
db.Close()
return nil, fmt.Errorf("连接数据库 %s 失败: %w", path, err)
}
return db, nil
}
// migrations 按顺序存放每一版的迁移语句。
//
// 外层一个元素 = 一个版本;内层是该版本要执行的语句,**一条一执行**。
// 不把多条语句塞进一个字符串,是因为 database/sql 的 Exec 对
// "一次执行多条语句"的支持因驱动而异,拆开最稳妥。
//
// 加新版本时**只能往末尾追加**,不许改动已有元素——
// 已经发布出去的库是按旧语句建的,改了会导致新旧库结构不一致。
//
// [必须] 这条规则本身也写进了 admin/AGENTS.md:#20 就是因为 v1 被原地改写、
// 而老库的 user_version 已经越过了它,那次改写永远不会在老库上重跑,
// 程序拿着一个和代码对不上的库静默启动。v3 不在这个 slice 里
// (见下面 schemaVersion 和 migrateV3 的注释),也是同一个教训的直接结果:
// v3 要做的事超出"一串 SQL 顺序执行",硬塞进这个只支持纯 SQL 的结构反而
// 掩盖了它的特殊性。
var migrations = [][]string{
// v1: 初始表结构,对应 docs/admin/03-data-model.md
{
`CREATE TABLE shopee_products (
goods_id TEXT PRIMARY KEY,
title TEXT NOT NULL,
shopee_status TEXT,
main_sku_code TEXT,
-- 下面三个是人工维护的,报表里没有,导入时绝不能覆盖
pdd_goods_url TEXT,
pdd_goods_id TEXT,
pdd_data TEXT,
collect_status TEXT NOT NULL DEFAULT 'no_link'
CHECK (collect_status IN (
'no_link', 'pending', 'collecting',
'collected', 'failed'
)),
collect_error TEXT,
collected_at TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`,
`CREATE INDEX idx_shopee_products_status ON shopee_products(collect_status);`,
`CREATE TABLE shopee_skus (
sku_id TEXT PRIMARY KEY,
goods_id TEXT NOT NULL,
spec_raw TEXT NOT NULL,
color TEXT,
size TEXT,
advice TEXT,
parse_ok INTEGER NOT NULL DEFAULT 0,
sku_code TEXT,
is_manual INTEGER NOT NULL DEFAULT 0,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
FOREIGN KEY (goods_id) REFERENCES shopee_products(goods_id) ON DELETE CASCADE
);`,
`CREATE INDEX idx_shopee_skus_goods ON shopee_skus(goods_id);`,
`CREATE INDEX idx_shopee_skus_parse ON shopee_skus(parse_ok);`,
`CREATE TABLE syb_orders (
syb_id TEXT PRIMARY KEY,
order_no TEXT NOT NULL,
title TEXT,
shopee_goods_id TEXT,
shopee_sku_id TEXT,
quantity INTEGER NOT NULL CHECK (quantity > 0),
price_twd_cent INTEGER CHECK (price_twd_cent IS NULL OR price_twd_cent >= 0),
image_url TEXT,
syb_data TEXT NOT NULL DEFAULT '{}',
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`,
`CREATE INDEX idx_syb_orders_order ON syb_orders(order_no);`,
`CREATE INDEX idx_syb_orders_goods ON syb_orders(shopee_goods_id);`,
`CREATE INDEX idx_syb_orders_list ON syb_orders(updated_at DESC, syb_id DESC);`,
`CREATE TABLE sku_mappings (
shopee_sku_id TEXT PRIMARY KEY,
goods_id TEXT NOT NULL,
pdd_options TEXT NOT NULL,
mapped_at TEXT NOT NULL,
mapped_by TEXT,
FOREIGN KEY (shopee_sku_id) REFERENCES shopee_skus(sku_id) ON DELETE CASCADE
);`,
`CREATE INDEX idx_sku_mappings_goods ON sku_mappings(goods_id);`,
`CREATE TABLE tasks (
task_id TEXT PRIMARY KEY,
task_type TEXT NOT NULL CHECK (task_type IN ('collect', 'purchase')),
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'assigned', 'claimed',
'succeeded', 'manual_review',
'failed', 'cancelled')),
version INTEGER NOT NULL DEFAULT 1 CHECK (version > 0),
priority INTEGER NOT NULL DEFAULT 0,
assigned_client TEXT,
claimed_at TEXT,
syb_id TEXT,
order_no TEXT,
goods_id TEXT,
shopee_sku_id TEXT,
-- Client 契约要求:pdd_goods_url 必填;
-- 采购任务的 quantity 和 max_price_cent 也必填(价格保护)
pdd_goods_url TEXT NOT NULL,
pdd_goods_id TEXT,
pdd_options TEXT,
quantity INTEGER CHECK (quantity IS NULL OR quantity > 0),
max_price_cent INTEGER CHECK (max_price_cent IS NULL OR max_price_cent > 0),
result_data TEXT,
error_code TEXT,
error_message TEXT,
finished_at TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`,
`CREATE INDEX idx_tasks_claim ON tasks(assigned_client, status, priority DESC, created_at);`,
`CREATE INDEX idx_tasks_list ON tasks(updated_at DESC, task_id DESC);`,
`CREATE INDEX idx_tasks_order ON tasks(order_no);`,
`CREATE TABLE clients (
client_id TEXT PRIMARY KEY,
name TEXT,
device_address TEXT,
platform TEXT,
pdd_package TEXT,
capabilities TEXT,
last_seen_at TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`,
`CREATE TABLE idempotency_keys (
key TEXT PRIMARY KEY,
request_hash TEXT NOT NULL,
response_body TEXT NOT NULL,
created_at TEXT NOT NULL
);`,
},
// v2: 领取历史。
//
// 为什么需要它:契约要求"只有**从未分配给该客户端**的任务才返回 403"
// (docs/admin/04-client-api.md §4.1)。但 tasks.assigned_client 只记
// **当前**归属,任务一旦重派给别人,就查不出原来那台领过——
// 而契约又明确要求"已重派仍要接受原客户端提交的结果"。
// 没有这张表,那条规则根本没法判断。
//
// 顺带得到一份审计记录:这个任务被哪几台客户端领过。
{
`CREATE TABLE task_claims (
task_id TEXT NOT NULL,
client_id TEXT NOT NULL,
claimed_at TEXT NOT NULL,
PRIMARY KEY (task_id, client_id)
);`,
`CREATE INDEX idx_task_claims_client ON task_claims(client_id);`,
},
}
// schemaVersion 是当前代码支持的最新 user_version。
//
// 之所以不是 len(migrations),是因为 v3(把 PDD 采集结果从 shopee_products
// 拆到独立的 pdd_products、重建 sku_mappings 主键)做的事超出了
// "一串 SQL 顺序执行":它必须在事务外切换 PRAGMA foreign_keys、
// 要在丢弃旧 sku_mappings 前用 Go 数出行数打日志。
// 这些事纯 SQL 表达不了,所以 v3 单独用 migrateV3 函数实现,
// 不放进 migrations 这个只支持"一条一条执行 SQL"的结构里。
//
// 背景见 #20:v1 曾经被原地改写而不是新增版本,导致已经建过库的机器
// (user_version 已经越过 v1)永远不会重跑改写后的语句,程序拿着一个
// 和代码对不上的库静默启动。
const schemaVersion = 8
// migrationV4 给 PDD 商品增加店铺名。
//
// v3 是特殊的 Go 迁移,不能塞进上面的纯 SQL migrations。v4 必须等 v3
// 建好 pdd_products 后再执行,所以单独放在这里。已经发布的 v1/v2 原文
// 保持不动,老库才能可靠地逐版升级。
var migrationV4 = []string{
`ALTER TABLE pdd_products ADD COLUMN shop_name TEXT;`,
}
// migrationV5 是工单 #46(顺运宝货运单同步)需要的三样东西:
// 会话缓存表、同步进度表、给 syb_orders 补一列规格原文。
//
// `[必须]` 三条都是新增(新表或 ADD COLUMN),不改动任何已发布的列/表,
// 见 admin/AGENTS.md「迁移只追加」。
var migrationV5 = []string{
// 会话缓存。只存 Cookie 就够——08 §3.1 已确认 JWT 从不参与后续请求认证,
// 缓存 token 没有意义,缓存 Cookie 才能免登录。
`CREATE TABLE syb_session (
username TEXT PRIMARY KEY,
cookies TEXT NOT NULL, -- JSON 数组
expires_at TEXT NOT NULL, -- min(JWT exp, 24h)
updated_at TEXT NOT NULL
);`,
// 上次同步到哪。只允许一行(CHECK id = 1),increment 同步靠它算日期范围。
`CREATE TABLE syb_sync_state (
id INTEGER PRIMARY KEY CHECK (id = 1),
last_synced_at TEXT,
updated_at TEXT NOT NULL
);`,
// 规格原文。匹配要用它,放列里才能查、才能在界面显示,
// 对应 shopee_skus.spec_raw,两边同一个概念。
`ALTER TABLE syb_orders ADD COLUMN product_spec TEXT;`,
}
// migrationV6 增加 Admin 网页账号和 Session。只追加新表和索引,
// v1-v5 的任何语句都不修改,见工单 #50。
var migrationV6 = []string{
`CREATE TABLE users (
user_id TEXT PRIMARY KEY,
username TEXT NOT NULL COLLATE NOCASE UNIQUE,
password_hash TEXT NOT NULL,
role TEXT NOT NULL CHECK (role IN ('admin', 'purchaser')),
status TEXT NOT NULL CHECK (status IN ('active', 'disabled')),
last_login_at TEXT,
password_changed_at TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`,
`CREATE TABLE web_sessions (
session_hash TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
expires_at TEXT NOT NULL,
created_at TEXT NOT NULL,
last_seen_at TEXT NOT NULL,
FOREIGN KEY (user_id) REFERENCES users(user_id)
);`,
`CREATE INDEX idx_web_sessions_user ON web_sessions(user_id);`,
`CREATE INDEX idx_web_sessions_expiry ON web_sessions(expires_at);`,
}
// migrationV7 记录客户端当前负责人及完整转交历史,见工单 #54。
// client_id 故意不加 clients 外键:客户端记录可删除后由同一稳定编号重新登记,
// 归属和审计历史不能因此丢失。
var migrationV7 = []string{
`CREATE TABLE client_user_assignments (
assignment_id TEXT PRIMARY KEY,
client_id TEXT NOT NULL,
user_id TEXT NOT NULL,
started_at TEXT NOT NULL,
ended_at TEXT,
assigned_by_user_id TEXT NOT NULL,
ended_by_user_id TEXT,
end_reason TEXT CHECK (end_reason IS NULL OR end_reason IN ('unbind', 'transfer')),
FOREIGN KEY (user_id) REFERENCES users(user_id),
FOREIGN KEY (assigned_by_user_id) REFERENCES users(user_id),
FOREIGN KEY (ended_by_user_id) REFERENCES users(user_id),
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))
);`,
`CREATE UNIQUE INDEX idx_client_assignment_current
ON client_user_assignments(client_id) WHERE ended_at IS NULL;`,
`CREATE INDEX idx_client_assignment_user
ON client_user_assignments(user_id, ended_at, client_id);`,
`CREATE INDEX idx_client_assignment_history
ON client_user_assignments(client_id, started_at DESC);`,
}
// migrationV8 持久化顺运宝同步记录,见工单 #59。
// 只新增表和索引,不修改 v1-v7 的任何已发布语句。
var migrationV8 = []string{
`CREATE TABLE syb_sync_runs (
run_id TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
date_from TEXT NOT NULL,
date_to TEXT NOT NULL,
status TEXT NOT NULL CHECK (status IN ('running', 'succeeded', 'failed', 'interrupted')),
stock_count INTEGER NOT NULL DEFAULT 0,
detail_count INTEGER NOT NULL DEFAULT 0,
created_count INTEGER NOT NULL DEFAULT 0,
updated_count INTEGER NOT NULL DEFAULT 0,
skipped_count INTEGER NOT NULL DEFAULT 0,
error_message TEXT,
cursor_advanced INTEGER NOT NULL DEFAULT 0 CHECK (cursor_advanced IN (0, 1)),
started_at TEXT NOT NULL,
finished_at TEXT,
FOREIGN KEY (user_id) REFERENCES users(user_id),
CHECK ((status = 'running' AND finished_at IS NULL)
OR (status <> 'running' AND finished_at IS NOT NULL))
);`,
`CREATE INDEX idx_syb_sync_runs_started
ON syb_sync_runs(started_at DESC, run_id DESC);`,
}
// Migrate 把数据库升到最新版本。
// 已经是最新的就什么都不做,可以重复调用。
func Migrate(db *sql.DB) error {
var current int
if err := db.QueryRow("PRAGMA user_version").Scan(&current); err != nil {
return fmt.Errorf("读取 user_version 失败: %w", err)
}
if current > schemaVersion {
return fmt.Errorf(
"数据库版本 %d 高于本程序支持的 %d,"+
"说明这个库是更新版本的程序建的,请升级程序而不是降级",
current, schemaVersion)
}
reached := current
for v := current; v < len(migrations); v++ {
tx, err := db.Begin()
if err != nil {
return fmt.Errorf("开始迁移 v%d 失败: %w", v+1, err)
}
for i, stmt := range migrations[v] {
if _, err := tx.Exec(stmt); err != nil {
tx.Rollback()
return fmt.Errorf("执行迁移 v%d 第 %d 条语句失败: %w", v+1, i+1, err)
}
}
// PRAGMA 不支持参数化,这里的值来自循环变量而非外部输入,安全。
if _, err := tx.Exec(fmt.Sprintf("PRAGMA user_version = %d", v+1)); err != nil {
tx.Rollback()
return fmt.Errorf("更新 user_version 到 %d 失败: %w", v+1, err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("提交迁移 v%d 失败: %w", v+1, err)
}
reached = v + 1
}
// v3:见 schemaVersion 的注释,为什么它不在 migrations 里、
// 单独用一个函数处理。全新库也会先走完 v1/v2(拿到旧版 shopee_products /
// sku_mappings 结构),再由这一步收敛成最终结构——
// 这样"全新库"和"老库升级"最终跑的是完全相同的 v3 代码,
// 不需要分别维护两条路径。
if reached < 3 {
if err := migrateV3(db); err != nil {
return err
}
reached = 3
}
// v4 是普通的追加列迁移,但必须排在特殊 v3 后面执行。
if reached < 4 {
if err := runSQLMigration(db, 4, migrationV4); err != nil {
return err
}
}
// v5 同理,是普通的建表 + 追加列迁移,排在 v4 后面执行。
if reached < 5 {
if err := runSQLMigration(db, 5, migrationV5); err != nil {
return err
}
reached = 5
}
// v6 是纯追加的用户和 Web Session 表。
if reached < 6 {
if err := runSQLMigration(db, 6, migrationV6); err != nil {
return err
}
reached = 6
}
// v7 是纯追加的客户端归属历史表和索引。
if reached < 7 {
if err := runSQLMigration(db, 7, migrationV7); err != nil {
return err
}
reached = 7
}
// v8 只新增顺运宝同步记录表和倒序索引。
if reached < 8 {
if err := runSQLMigration(db, 8, migrationV8); err != nil {
return err
}
}
return nil
}
// runSQLMigration 在一个事务里执行指定版本的 SQL,并最后更新 user_version。
// 它只接收本文件中写死的版本号和 SQL,不接收外部输入。
func runSQLMigration(db *sql.DB, version int, statements []string) error {
tx, err := db.Begin()
if err != nil {
return fmt.Errorf("开始迁移 v%d 失败: %w", version, err)
}
defer tx.Rollback()
for i, stmt := range statements {
if _, err := tx.Exec(stmt); err != nil {
return fmt.Errorf("执行迁移 v%d 第 %d 条语句失败: %w", version, i+1, err)
}
}
if _, err := tx.Exec(fmt.Sprintf("PRAGMA user_version = %d", version)); err != nil {
return fmt.Errorf("更新 user_version 到 %d 失败: %w", version, err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("提交迁移 v%d 失败: %w", version, err)
}
return nil
}
// migrateV3 把库收敛成当前结构,对应工单 #20。
//
// # 起点不止一种,不能用 user_version 推断结构
//
// #16 曾经原地改写过 v1(把 pdd_products 等结构直接塞进 v1,没有新增版本号),
// 所以 user_version = 2 现在对应两种不同的**真实**结构:
//
// 原始 v1 建的库 没有 pdd_products;shopee_products 还带着
// pdd_data / collect_status 等四列;sku_mappings
// 还是单列主键。
// #16 改写后的 v1 建的库 结构已经是最终形态,只是 user_version 还停在 2。
//
// 只要用 user_version 是不是 2 来判断"要不要迁移"就会踩空——第二种库
// 一旦被当成第一种处理,"建 pdd_products" 这一步会直接撞 "table already
// exists"。所以下面每一步都必须先查**实际结构**(sqlite_master /
// PRAGMA table_info),再决定做不做;已经是最终形态的库,这个函数应该
// 什么都不改,只把版本号推到 3。
//
// # 四件事,各自独立判断要不要做
//
// 1. 建 pdd_products 和索引——表已存在就跳过。
// 2. 把 shopee_products 上的 PDD 采集数据搬过去,按 pdd_goods_id 去重——
// shopee_products 已经没有 pdd_data 列就跳过(数据要么搬过了,
// 要么这张表从来就没有过)。
// 3. 重建 sku_mappings——已经有 pdd_goods_id 列(说明已经是新结构)就跳过;
// 否则旧数据全部丢弃(新主键需要的 pdd_option_key 是
// service.OptionKey() 用 json.Marshal 算出来的,SQL 复现不了,硬凑
// 有静默买错东西的风险,见 admin/AGENTS.md),丢之前先数出行数打日志。
// 4. 重建 shopee_products——还带着旧的四列中任意一个就重建,去掉那四列;
// 四列都已经不在就跳过。
//
// 重建 shopee_products 时必须按 SQLite 官方的表重建流程处理外键
// (https://www.sqlite.org/lang_altertable.html 的 12 步):
// shopee_skus.goods_id 有外键指向 shopee_products(goods_id) ON DELETE CASCADE,
// 如果不先关闭外键检查就 DROP/RENAME 这张表,可能把 shopee_skus 的数据带跑。
// 而 PRAGMA foreign_keys 只能在没有打开事务时切换,所以必须先关、
// 再开事务,事务里全部做完再提交、最后恢复。这套开销只在真的要重建
// shopee_products 时才需要——sku_mappings 只是外键的子表(没有别的表
// 指向它),丢它不会牵连别的表,不需要关闭外键检查。
func migrateV3(db *sql.DB) error {
const toVersion = 3
ctx := context.Background()
needCreatePddProducts, err := tableMissing(db, "pdd_products")
if err != nil {
return fmt.Errorf("迁移 v3 检查 pdd_products 是否存在失败: %w", err)
}
shopeeCols, err := tableColumnSet(db, "shopee_products")
if err != nil {
return fmt.Errorf("迁移 v3 检查 shopee_products 结构失败: %w", err)
}
needDedup := shopeeCols["pdd_data"]
needRebuildShopeeProducts := shopeeCols["pdd_data"] || shopeeCols["collect_status"] ||
shopeeCols["collect_error"] || shopeeCols["collected_at"]
skuCols, err := tableColumnSet(db, "sku_mappings")
if err != nil {
return fmt.Errorf("迁移 v3 检查 sku_mappings 结构失败: %w", err)
}
needRebuildSkuMappings := !skuCols["pdd_goods_id"]
if !needCreatePddProducts && !needDedup && !needRebuildShopeeProducts && !needRebuildSkuMappings {
// 结构已经是最终形态(#16 改写后的 v1 建的库),什么都不用改,
// 只需要把版本号推到 3。不打"新增 / 重建"那行日志——那是假话,
// 会误导操作员以为数据被动过。
log.Printf("数据库迁移 v3:结构已是最新,仅更新版本号")
if _, err := db.Exec(fmt.Sprintf("PRAGMA user_version = %d", toVersion)); err != nil {
return fmt.Errorf("更新 user_version 到 %d 失败: %w", toVersion, err)
}
return nil
}
// 按**实际做了什么**打日志,不是无条件打同一行——操作员/维护者要能
// 从启动日志一眼看出发生过一次结构迁移,不用等到点坏页面才知道库被动过,
// 这正是 #20 要修的"静默"问题;但日志内容必须真实,不能不管做没做都
// 打同一句。
var actions []string
if needCreatePddProducts {
actions = append(actions, "新增 pdd_products")
}
if needRebuildShopeeProducts {
actions = append(actions, "重建 shopee_products")
}
if needRebuildSkuMappings {
actions = append(actions, "重建 sku_mappings")
}
log.Printf("数据库迁移 v3:%s", strings.Join(actions, "、"))
var conn *sql.Conn
if needRebuildShopeeProducts {
// 只有要重建 shopee_products 才需要这套连接和 PRAGMA 切换——见函数
// 顶部的注释。单独拿一条连接:PRAGMA foreign_keys=OFF 之后紧接着要在
// **同一条连接**上开事务,database/sql 的连接池不保证 db.Exec 和
// db.Begin 用的是同一条连接。
conn, err = db.Conn(ctx)
if err != nil {
return fmt.Errorf("迁移 v3 获取专用连接失败: %w", err)
}
defer conn.Close()
if _, err := conn.ExecContext(ctx, "PRAGMA foreign_keys=OFF"); err != nil {
return fmt.Errorf("迁移 v3 关闭外键检查失败: %w", err)
}
// 这条连接用完会回到连接池,后面别的代码还会拿到它继续用。
// 不管上面成功还是失败都要把外键检查恢复成 ON——
// Open() 承诺过整个连接池的外键检查是打开的,这里关了就要负责关回去。
// 这个 defer 注册在 conn.Close() 之后,按 LIFO 顺序会先于 Close 执行。
defer func() {
if _, err := conn.ExecContext(context.Background(), "PRAGMA foreign_keys=ON"); err != nil {
log.Printf("警告: 迁移 v3 后恢复外键检查失败: %v", err)
}
}()
}
var tx *sql.Tx
if conn != nil {
tx, err = conn.BeginTx(ctx, nil)
} else {
tx, err = db.BeginTx(ctx, nil)
}
if err != nil {
return fmt.Errorf("开始迁移 v3 事务失败: %w", err)
}
defer tx.Rollback() // 已提交的事务再 Rollback 是空操作,安全
// ① 建 pdd_products 和索引。内容照搬当前 v1 里的定义,含全部注释——
// 这是最终要收敛到的表结构,不能和 v1 曾经的定义有任何出入。
if needCreatePddProducts {
for i, stmt := range migrateV3CreatePddProducts {
if _, err := tx.ExecContext(ctx, stmt); err != nil {
return fmt.Errorf("迁移 v3 建 pdd_products 第 %d 条语句失败: %w", i+1, err)
}
}
}
// ② 把老 shopee_products 上的 PDD 数据搬过去,按 pdd_goods_id 去重。
// 必须在下面重建 shopee_products 之前执行——一旦 shopee_products 被重建,
// pdd_data / collect_status / collect_error / collected_at 这几列就没了。
if needDedup {
if _, err := tx.ExecContext(ctx, migrateV3DedupIntoPddProducts); err != nil {
return fmt.Errorf("迁移 v3 搬运 PDD 数据失败: %w", err)
}
}
// ③ 旧 sku_mappings 全部丢弃:新主键需要的 pdd_option_key 是
// service.OptionKey() 用 json.Marshal 算出来的,SQL 复现不了;
// 硬凑一个键有静默买错东西的风险,见 admin/AGENTS.md。
// 丢之前先数出行数打日志,不能悄悄丢——N 为 0 时不打,
// 避免每次启动都刷一行没用的日志。
if needRebuildSkuMappings {
var discardedMappings int
if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM sku_mappings`).Scan(&discardedMappings); err != nil {
return fmt.Errorf("迁移 v3 统计旧 sku_mappings 行数失败: %w", err)
}
if discardedMappings > 0 {
log.Printf("迁移 v3:丢弃了 %d 条旧规格映射(缺少 pdd_goods_id / pdd_option_key,无法安全迁移,请重新匹配)", discardedMappings)
}
for i, stmt := range migrateV3RebuildSkuMappings {
if _, err := tx.ExecContext(ctx, stmt); err != nil {
return fmt.Errorf("迁移 v3 重建 sku_mappings 第 %d 条语句失败: %w", i+1, err)
}
}
}
// ④ 重建 shopee_products:去掉已经搬去 pdd_products 的四个列。
// 按 SQLite 12 步流程:建新表 -> 搬数据 -> 删旧表 -> 改名 -> 建索引。
if needRebuildShopeeProducts {
for i, stmt := range migrateV3RebuildShopeeProducts {
if _, err := tx.ExecContext(ctx, stmt); err != nil {
return fmt.Errorf("迁移 v3 重建 shopee_products 第 %d 条语句失败: %w", i+1, err)
}
}
// 外键检查原本是开着的:重建完必须确认没有把 shopee_skus 的数据带丢
// (比如误伤了它和 shopee_products 之间的外键关系)。
if err := checkForeignKeys(ctx, tx); err != nil {
return fmt.Errorf("迁移 v3 外键校验失败: %w", err)
}
}
if _, err := tx.ExecContext(ctx, fmt.Sprintf("PRAGMA user_version = %d", toVersion)); err != nil {
return fmt.Errorf("更新 user_version 到 %d 失败: %w", toVersion, err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("提交迁移 v3 失败: %w", err)
}
return nil
}
// tableMissing 判断某张表是否不存在。
func tableMissing(db *sql.DB, table string) (bool, error) {
var name string
err := db.QueryRow(
`SELECT name FROM sqlite_master WHERE type = 'table' AND name = ?`, table,
).Scan(&name)
if err == sql.ErrNoRows {
return true, nil
}
if err != nil {
return false, err
}
return false, nil
}
// tableColumnSet 返回某张表当前的列名集合,供 migrateV3 判断
// "这张表是不是旧结构" 用——user_version 推断不出真实结构,见 migrateV3 顶部注释。
func tableColumnSet(db *sql.DB, table string) (map[string]bool, error) {
// table 只来自本文件里写死的表名常量,不是外部输入,字符串拼接是安全的
// (PRAGMA 本身也不支持参数化,PRAGMA user_version 也是这么处理的)。
rows, err := db.Query(`PRAGMA table_info(` + table + `)`)
if err != nil {
return nil, err
}
defer rows.Close()
cols := map[string]bool{}
for rows.Next() {
var cid, notnull, pk int
var name, ctype string
var dflt sql.NullString
if err := rows.Scan(&cid, &name, &ctype, &notnull, &dflt, &pk); err != nil {
return nil, err
}
cols[name] = true
}
return cols, rows.Err()
}
// migrateV3CreatePddProducts 建 pdd_products 和索引。
//
// [必须] 这段 SQL 文本必须和 `git show 998c06a:admin/repository/db.go`
// 里 pdd_products 的定义逐字节一致(含缩进),不能只是"看起来一样"。
//
// 原因:#16(998c06a)把这张表直接建进了 v1,那批库上 sqlite_master 存的
// 就是 998c06a 里这段文本原样的字节(连同当时 gofmt 缩进出来的那个前导
// TAB)。这个函数只在表**不存在**时才会执行这条 CREATE(见 migrateV3 里
// needCreatePddProducts 的判断),所以已经建过表的库不会被这段文本重新
// 覆盖——"新建库" 和 "已经带 pdd_products 的老库" 要收敛到完全相同的
// sqlite_master.sql,就必须让新建的这份文本和老库上躺着的那份逐字节相同,
// 哪怕只差一个空格/TAB 都会让 TestMigrate_不同起点最终schema一致 失败
// (这个坑已经在 #20 审查阶段被变异测试连同真实文本对比一起抓到过一次)。
var migrateV3CreatePddProducts = []string{
`CREATE TABLE pdd_products (
id INTEGER PRIMARY KEY AUTOINCREMENT,
-- 从 PDD 链接里解析出来。它不是主键,所以**必须加 UNIQUE**:
-- 少了这条约束,同一个 PDD 商品会被存成好几行,
-- 采好几遍,映射还说不清指向哪一行。
goods_id TEXT NOT NULL UNIQUE,
url TEXT NOT NULL, -- 操作员填的链接原文
title TEXT, -- 采集回来,人工核对"是不是我要的那个商品"
skus_json TEXT, -- schema_version + dimensions + skus
-- 注意这里**没有 no_link**:这张表里有这一行,就说明链接已经填了。
-- "未填链接"是蝦皮侧的状态(shopee_products.pdd_goods_id 为空)。
collect_status TEXT NOT NULL DEFAULT 'pending'
CHECK (collect_status IN (
'pending', 'collecting', 'collected', 'failed'
)),
collect_msg TEXT, -- 失败原因,要能定位问题
artifact_ref TEXT, -- 诊断产物在哪台机器哪个目录
collected_at TEXT,
-- 软删除。不硬删是因为 sku_mappings 指向它,
-- 硬删会把人工攒了很久的匹配成果一起带走。
deleted_at TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`,
`CREATE INDEX idx_pdd_products_status ON pdd_products(collect_status);`,
}
// migrateV3DedupIntoPddProducts 按 pdd_goods_id 去重,把老
// shopee_products 上人工攒的 PDD 数据搬进 pdd_products。
//
// - MIN(COALESCE(pdd_goods_url, 空字符串)) / MIN(created_at):多行里任取一个即可,
// 用 MIN 只是为了确定性(同一批输入每次跑结果一样,方便排查)。
// - MAX(pdd_data) / MAX(collect_error) / MAX(collected_at):同理,
// 只是要"取到某一行的值",用 MAX 是同一个考虑。
// - CASE MAX(collect_status) ...:collect_status 只有全组都是
// 'collected' 时才判定为 collected('collected' 按字符串比较是这几个
// 取值里最小的,只要组里有任何一行不是 collected,MAX 就会取到别的值);
// 只要没有 collected 但有 failed 就判定 failed;
// 其余(含 collecting / no_link / pending)一律落进 ELSE,判定 pending
// —— 这正好同时满足"collecting 映射成 pending"和
// "no_link 映射成 pending,不撞新 CHECK(新表里没有 no_link)"两条要求。
var migrateV3DedupIntoPddProducts = `
INSERT INTO pdd_products
(goods_id, url, title, skus_json, collect_status,
collect_msg, collected_at, created_at, updated_at)
SELECT
pdd_goods_id,
MIN(COALESCE(pdd_goods_url, '')),
NULL,
MAX(pdd_data),
CASE MAX(collect_status)
WHEN 'collected' THEN 'collected'
WHEN 'failed' THEN 'failed'
ELSE 'pending'
END,
MAX(collect_error),
MAX(collected_at),
MIN(created_at), MAX(updated_at)
FROM shopee_products
WHERE pdd_goods_id IS NOT NULL AND pdd_goods_id <> ''
GROUP BY pdd_goods_id;`
// migrateV3RebuildShopeeProducts 重建 shopee_products:
// 去掉已经搬去 pdd_products 的 pdd_data / collect_status / collect_error /
// collected_at 四列,删掉跟着它们的 idx_shopee_products_status,
// 换成新结构需要的 idx_shopee_products_pdd。
//
// # 步骤顺序为什么是"建 _new → 搬数据 → 删旧表 → 改名"(标准 12 步顺序)
//
// 这里**必须**按 SQLite 官方 12 步流程的顺序来,不能改成"先把旧表改名
// 让开、再直接用最终表名建新表"(表面上能避开下面说的引号问题,
// 实测过、但会坏得更彻底):
//
// shopee_skus.goods_id 有 `FOREIGN KEY (goods_id) REFERENCES
// shopee_products(goods_id)`。SQLite 的 `ALTER TABLE ... RENAME TO`
// **默认会连带更新别的表里引用这张表的外键定义**——如果改成先把
// shopee_products RENAME 成 shopee_products_old,shopee_skus 的外键子句
// 会被自动重写成 `REFERENCES "shopee_products_old"(goods_id)`;
// 后面一 DROP TABLE shopee_products_old,shopee_skus 就带着一条指向
// 不存在的表的外键,`PRAGMA foreign_key_check` 直接报错,
// 而且这个坏结果比"多一对引号"严重得多。
//
// 按标准顺序(建 shopee_products_new → 搬数据 → DROP 掉的是旧的
// shopee_products,不是被引用的名字 → RENAME shopee_products_new
// 成 shopee_products)不会触发这个重写:shopee_skus 的外键子句
// 全程写的都是 "shopee_products" 这个名字,没有变过,RENAME 结束后
// 这个名字重新存在,直接就能解析上,不需要 SQLite 帮它改写任何东西。
// 实测两种顺序的效果:
//
// 方案:先建 _new 再 RENAME 成正式名(本文件采用的顺序)
// shopee_skus 的外键子句:REFERENCES shopee_products(goods_id) 不变 ✓
// shopee_products 自己的 CREATE TABLE 文本:被 RENAME 加上引号 CREATE TABLE "shopee_products" (...)
//
// 方案:先把旧表 RENAME 让开,再直接用正式名建新表
// shopee_skus 的外键子句:被自动重写成 REFERENCES "shopee_products_old"(goods_id) ← 引用悬空,PRAGMA foreign_key_check 报错
//
// 所以选择"忍受 shopee_products 自己的 CREATE TABLE 文本被套一层引号",
// 而不是"外键指向一张已经被删掉的表"。引号这个副作用改在
// assertSchemaEqual 的比较逻辑里通过归一化解决(见 migrate_test.go
// 里 normalizeCreateStatement 的注释),不在这里想办法回避——
// 这个坑是 #20 审查阶段三起点收敛测试跑出来才发现的,光读代码看不出来。
//
// [必须] 除了表名后缀 _new、以及 RENAME 带来的引号,这段 SQL 里
// shopee_products 的列定义必须和
// `git show 998c06a:admin/repository/db.go` 里的定义逐字节一致
// (含缩进),理由同 migrateV3CreatePddProducts 的注释。
var migrateV3RebuildShopeeProducts = []string{
`CREATE TABLE shopee_products_new (
goods_id TEXT PRIMARY KEY,
title TEXT NOT NULL,
shopee_status TEXT,
main_sku_code TEXT,
-- 人工维护的,蝦皮报表里没有这两列,Excel 导入时绝不能覆盖。
-- pdd_goods_id 指向 pdd_products.goods_id,表示"这个蝦皮商品
-- 当前对应哪个 PDD 商品"。PDD 商品下架换新时改这里。
pdd_goods_url TEXT,
pdd_goods_id TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`,
`INSERT INTO shopee_products_new
(goods_id, title, shopee_status, main_sku_code,
pdd_goods_url, pdd_goods_id, created_at, updated_at)
SELECT
goods_id, title, shopee_status, main_sku_code,
pdd_goods_url, pdd_goods_id, created_at, updated_at
FROM shopee_products;`,
`DROP TABLE shopee_products;`,
`ALTER TABLE shopee_products_new RENAME TO shopee_products;`,
`CREATE INDEX idx_shopee_products_pdd ON shopee_products(pdd_goods_id);`,
}
// migrateV3RebuildSkuMappings 重建 sku_mappings。
// 旧数据全部丢弃(调用方已经数过行数打过日志),不搬任何数据——
// 新主键需要的 pdd_option_key 只有 Go 的 service.OptionKey() 能算,
// SQL 里凑不出来,硬凑有静默买错东西的风险。
//
// 这里是直接 DROP 旧表、CREATE 新表(同一个最终表名,不经过 RENAME),
// 不会有 migrateV3RebuildShopeeProducts 注释里说的引号问题。
//
// [必须] 这段 SQL 文本必须和 `git show 998c06a:admin/repository/db.go`
// 里 sku_mappings 的定义逐字节一致(含缩进),理由同
// migrateV3CreatePddProducts 的注释。
var migrateV3RebuildSkuMappings = []string{
`DROP TABLE sku_mappings;`,
`CREATE TABLE sku_mappings (
shopee_sku_id TEXT NOT NULL,
pdd_goods_id TEXT NOT NULL,
pdd_option_key TEXT NOT NULL,
pdd_options TEXT NOT NULL,
goods_id TEXT NOT NULL,
mapped_at TEXT NOT NULL,
mapped_by TEXT,
PRIMARY KEY (shopee_sku_id, pdd_goods_id),
FOREIGN KEY (shopee_sku_id) REFERENCES shopee_skus(sku_id) ON DELETE CASCADE
);`,
`CREATE INDEX idx_sku_mappings_goods ON sku_mappings(goods_id);`,
`CREATE INDEX idx_sku_mappings_pdd ON sku_mappings(pdd_goods_id);`,
}
// checkForeignKeys 跑 PRAGMA foreign_key_check,有任何一行结果
// 就说明外键关系被破坏了(比如子表指向了一个已经不存在的父行)。
func checkForeignKeys(ctx context.Context, tx *sql.Tx) error {
rows, err := tx.QueryContext(ctx, "PRAGMA foreign_key_check")
if err != nil {
return fmt.Errorf("执行外键校验失败: %w", err)
}
defer rows.Close()
if rows.Next() {
return fmt.Errorf("存在外键约束冲突,重建表的过程把关联数据带丢了")
}
return rows.Err()
}
// requiredTables 是当前代码依赖的全部表。
// Migrate 跑完之后用它做一次自检,见 CheckSchema。
var requiredTables = []string{
"shopee_products", "shopee_skus", "pdd_products",
"syb_orders", "sku_mappings", "tasks", "clients",
"idempotency_keys", "task_claims",
"syb_session", "syb_sync_state",
"users", "web_sessions", "client_user_assignments", "syb_sync_runs",
}
// requiredColumns 只列出不能靠“表存在”发现的关键追加列。
// shop_name 是 v4 新增列;product_spec 是 v5 新增列——两者缺失时
// 查询对应页面会直接失败,因此启动时就应给出明确错误,
// 而不是等操作员点到页面才暴露。
var requiredColumns = map[string][]string{
"pdd_products": {"shop_name"},
"syb_orders": {"product_spec"},
"syb_sync_runs": {"user_id", "date_from", "date_to", "status", "cursor_advanced", "started_at", "finished_at"},
}
// CheckSchema 在 Migrate 成功后调用,确认代码依赖的表都在。
//
// [必须] 缺表就返回错误,调用方要**拒绝启动**,不是打个警告继续跑。
// #20 的教训就是静默启动:程序拿着一个和代码对不上的库正常起来了,
// 错误要等操作员点到那个页面才暴露——如果那是个写操作页面,
// 暴露出来的就不是报错而是写坏数据。
//
// 表名全部检查;关键的追加列也检查,防止 user_version 已更新但迁移未完整
// 落地时,程序拿着缺列的库继续启动。
func CheckSchema(db *sql.DB) error {
rows, err := db.Query(`SELECT name FROM sqlite_master WHERE type = 'table'`)
if err != nil {
return fmt.Errorf("读取数据库表清单失败: %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("读取数据库表清单失败: %w", err)
}
existing[name] = true
}
if err := rows.Err(); err != nil {
return fmt.Errorf("读取数据库表清单失败: %w", err)
}
var missing []string
for _, t := range requiredTables {
if !existing[t] {
missing = append(missing, t)
}
}
if len(missing) > 0 {
return fmt.Errorf(
"数据库结构与本程序不匹配:缺少表 %s。\n"+
"这通常是数据库比程序旧、而迁移没有覆盖到。\n"+
"请备份 data/admin.db 后删除它让程序重建,或联系维护者。",
strings.Join(missing, "、"))
}
for table, columns := range requiredColumns {
existingColumns, err := tableColumnSet(db, table)
if err != nil {
return fmt.Errorf("检查数据表 %s 的列失败: %w", table, err)
}
for _, column := range columns {
if !existingColumns[column] {
return fmt.Errorf(
"数据库结构与本程序不匹配:数据表 %s 缺少列 %s。\n"+
"这通常是数据库迁移没有完整执行,请备份 data/admin.db 后联系维护者。",
table, column)
}
}
}
return nil
}