Files
cmautobuy/admin/repository/db.go
T
chengmaandClaude Opus 5 998c06a2bf feat: PDD 商品数据独立成表 (#16)
原来 pdd_data 是 shopee_products 上的一个 JSON 字段,两个蝦皮商品指向
同一个 PDD 链接时会各存一份、各采一次;collect_status 描述的是 PDD 商品的
状态,却挂在蝦皮商品上,两份可能不一致。

更要紧的是 PDD 商品变动频繁(A 下架就得换 B),而 sku_mappings 只按
shopee_sku_id 做键——换商品后旧映射还在,B 恰好有同名规格但完全是另一件货
时会静默买错,事后查不出来。

改动
- 新增 pdd_products 表:id 主键 + goods_id UNIQUE + 4 个状态值(去掉
  no_link,「未填链接」改由 shopee_products.pdd_goods_id 为空表达)+
  软删除可复活
- shopee_products 去掉 pdd_data / collect_status / collect_error /
  collected_at,pdd_goods_id 改为引用
- sku_mappings 主键改为 (shopee_sku_id, pdd_goods_id),新增 pdd_option_key。
  查映射永远带上当前 PDD 商品,换商品后天然查不到旧映射,不需要删数据;
  换回原商品时旧映射直接复用
- 新增 OptionKey():用 json.Marshal 实现(Go 序列化 map 按键名排序,
  天然规范化),不自己拼字符串——规格文字里可能含 = 或 ;。
  存映射和查 SKU 必须用同一个函数,各写一遍会静默算出不同结果
- 采集结果改落 pdd_products,新增两条校验:
  返回的 goods_id 与请求不符 → 整体回滚拒绝(422),不静默存下;
  skus 为空数组 → 置 failed 而非 collected,否则界面显示"已采集"
  但数据毫无用处

实施时超出工单但必要的三处
- TaskExists 重构为 GetTaskInfo:原函数只返回蝦皮 goods_id,
  而采集结果要按 PDD goods_id 落库,不改取不到正确的键
- 复活时一并清空旧采集结果(skus_json / collect_msg / collected_at),
  否则复活后会显示"已采集"但数据是删除前的
- 删除 repository/shopee.go:两个函数签名全变且已迁到 pdd.go,留着是死代码

已验证(Go 1.23.0)
- go vet / gofmt / go test 全过,55 个测试
- 端到端补验了工单未覆盖的 HTTP 层:goods_id 不符返回 422
  COLLECT_GOODS_MISMATCH 且整体回滚(skus_json 空、任务仍 claimed、
  幂等记录 0 条);skus 为空返回 200 但状态 failed

遗留
- MarkCollecting / SoftDeletePddProduct 暂无调用方,等界面工单接上
- artifact_ref 存 diagnostics 原始 JSON,未按 client-001:artifacts/... 规范化,
  因 Client 侧尚未定义 diagnostics 结构
- 界面未实现(工单明确排除),四个页面仍为骨架

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-07 10:27:39 +08:00

310 lines
12 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 (
"database/sql"
"fmt"
"path/filepath"
// 纯 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 对
// "一次执行多条语句"的支持因驱动而异,拆开最稳妥。
//
// 加新版本时**只能往末尾追加**,不许改动已有元素——
// 已经发布出去的库是按旧语句建的,改了会导致新旧库结构不一致。
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,
-- 人工维护的,蝦皮报表里没有这两列,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
);`,
`CREATE INDEX idx_shopee_products_pdd ON shopee_products(pdd_goods_id);`,
`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 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);`,
`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 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);`,
`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);`,
},
}
// 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 > len(migrations) {
return fmt.Errorf(
"数据库版本 %d 高于本程序支持的 %d,"+
"说明这个库是更新版本的程序建的,请升级程序而不是降级",
current, len(migrations))
}
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)
}
}
return nil
}