feat: 保存并展示 PDD 店铺与价格采样信息 (#31)

This commit is contained in:
chengma
2026-08-07 17:13:25 +08:00
parent a2fc1501c4
commit 93beef8b36
17 changed files with 492 additions and 73 deletions
+72 -11
View File
@@ -264,7 +264,16 @@ var migrations = [][]string{
// 背景见 #20:v1 曾经被原地改写而不是新增版本,导致已经建过库的机器
// (user_version 已经越过 v1)永远不会重跑改写后的语句,程序拿着一个
// 和代码对不上的库静默启动。
const schemaVersion = 3
const schemaVersion = 4
// migrationV4 给 PDD 商品增加店铺名。
//
// v3 是特殊的 Go 迁移,不能塞进上面的纯 SQL migrations。v4 必须等 v3
// 建好 pdd_products 后再执行,所以单独放在这里。已经发布的 v1/v2 原文
// 保持不动,老库才能可靠地逐版升级。
var migrationV4 = []string{
`ALTER TABLE pdd_products ADD COLUMN shop_name TEXT;`,
}
// Migrate 把数据库升到最新版本。
// 已经是最新的就什么都不做,可以重复调用。
@@ -309,15 +318,46 @@ func Migrate(db *sql.DB) error {
// sku_mappings 结构),再由这一步收敛成最终结构——
// 这样"全新库"和"老库升级"最终跑的是完全相同的 v3 代码,
// 不需要分别维护两条路径。
if reached < schemaVersion {
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
}
}
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 推断结构
@@ -752,6 +792,13 @@ var requiredTables = []string{
"idempotency_keys", "task_claims",
}
// requiredColumns 只列出不能靠“表存在”发现的关键追加列。
// shop_name 是 v4 新增列;缺少它时查询 PDD 页面会直接失败,因此启动时
// 就应给出明确错误,而不是等操作员点到页面才暴露。
var requiredColumns = map[string][]string{
"pdd_products": {"shop_name"},
}
// CheckSchema 在 Migrate 成功后调用,确认代码依赖的表都在。
//
// [必须] 缺表就返回错误,调用方要**拒绝启动**,不是打个警告继续跑。
@@ -759,8 +806,8 @@ var requiredTables = []string{
// 错误要等操作员点到那个页面才暴露——如果那是个写操作页面,
// 暴露出来的就不是报错而是写坏数据。
//
// [建议] 只查表名,不逐列校验:够抓住"迁移没跑到、表没建出来"这一类问题,
// 代价也低。真出了列级别的不一致,业务 SQL 跑起来自然会报错。
// 表名全部检查;关键的追加列也检查,防止 user_version 已更新但迁移未完整
// 落地时,程序拿着缺列的库继续启动。
func CheckSchema(db *sql.DB) error {
rows, err := db.Query(`SELECT name FROM sqlite_master WHERE type = 'table'`)
if err != nil {
@@ -786,13 +833,27 @@ func CheckSchema(db *sql.DB) error {
missing = append(missing, t)
}
}
if len(missing) == 0 {
return nil
if len(missing) > 0 {
return fmt.Errorf(
"数据库结构与本程序不匹配:缺少表 %s。\n"+
"这通常是数据库比程序旧、而迁移没有覆盖到。\n"+
"请备份 data/admin.db 后删除它让程序重建,或联系维护者。",
strings.Join(missing, "、"))
}
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
}
+43
View File
@@ -452,6 +452,37 @@ func TestMigrate_已经是最新版本再次调用不报错(t *testing.T) {
}
}
func TestMigrate_v4增加可空店铺名且老数据保持NULL(t *testing.T) {
db := newV2NewStructureDB(t)
if _, err := db.Exec(`
INSERT INTO pdd_products (goods_id, url, collect_status, created_at, updated_at)
VALUES ('737116531267', 'https://example.invalid', 'pending', '2026-08-01T00:00:00Z', '2026-08-01T00:00:00Z')`); err != nil {
t.Fatalf("插入 v2 老数据失败: %v", err)
}
if err := Migrate(db); err != nil {
t.Fatalf("迁移到 v4 失败: %v", err)
}
columns, err := tableColumnSet(db, "pdd_products")
if err != nil {
t.Fatalf("读取 pdd_products 列失败: %v", err)
}
if !columns["shop_name"] {
t.Fatal("v4 必须新增 shop_name 列")
}
var shop sql.NullString
if err := db.QueryRow(
`SELECT shop_name FROM pdd_products WHERE goods_id = '737116531267'`,
).Scan(&shop); err != nil {
t.Fatalf("读取老数据店铺名失败: %v", err)
}
if shop.Valid {
t.Errorf("老数据的 shop_name 应为 NULL,实际 %q", shop.String)
}
}
// v2 新结构库(#16 改写后的 v1)第一次 Migrate 时,migrateV3 检测到结构
// 已经是最终形态,只更新版本号、不改任何表——第二次调用(对应用户重启
// 服务)应该是彻底的空操作:user_version 已经是 3,Migrate 最外层的
@@ -848,6 +879,18 @@ func TestCheckSchema_v2老库迁移后通过(t *testing.T) {
}
}
func TestCheckSchema_缺少v4关键列时拒绝(t *testing.T) {
db := newV2NewStructureDB(t)
if err := migrateV3(db); err != nil {
t.Fatalf("准备 v3 数据库失败: %v", err)
}
err := CheckSchema(db)
if err == nil || !strings.Contains(err.Error(), "shop_name") {
t.Fatalf("缺少 shop_name 时应拒绝启动并指出列名,实际 %v", err)
}
}
func TestCheckSchema_缺表时拒绝(t *testing.T) {
db := newFreshDB(t)
if _, err := db.Exec(`DROP TABLE pdd_products`); err != nil {
+13 -9
View File
@@ -12,7 +12,7 @@ import (
// pddColumns 是所有查询共用的列清单。
// 写成常量是为了让下面几个查询的列顺序和 scanPddProduct 永远对得上——
// 改列的时候只改这一处,改漏了会 Scan 到错的字段上,而且不报错。
const pddColumns = `id, goods_id, url, title, skus_json,
const pddColumns = `id, goods_id, url, title, shop_name, skus_json,
collect_status, collect_msg, artifact_ref, collected_at,
deleted_at, created_at, updated_at`
@@ -24,10 +24,10 @@ type rowScanner interface {
// scanPddProduct 按 pddColumns 的顺序读一行。
func scanPddProduct(s rowScanner, extra ...any) (*model.PddProduct, error) {
var p model.PddProduct
var title, skus, msg, artifact, collectedAt, deletedAt sql.NullString
var title, shopName, skus, msg, artifact, collectedAt, deletedAt sql.NullString
dest := []any{
&p.ID, &p.GoodsID, &p.URL, &title, &skus,
&p.ID, &p.GoodsID, &p.URL, &title, &shopName, &skus,
&p.CollectStatus, &msg, &artifact, &collectedAt,
&deletedAt, &p.CreatedAt, &p.UpdatedAt,
}
@@ -38,6 +38,7 @@ func scanPddProduct(s rowScanner, extra ...any) (*model.PddProduct, error) {
}
p.Title = title.String
p.ShopName = shopName.String
p.SkusJSON = skus.String
p.CollectMsg = msg.String
p.ArtifactRef = artifact.String
@@ -91,13 +92,13 @@ func EnsurePddProduct(q Execer, goodsID, url string) (*model.PddProduct, error)
// 复活:清掉删除标记,状态回到待采集。
// 采集结果一并清空——记录被删过一次,旧数据不能再当成有效的用。
//
// title 也要清:它是**采集回来的**(SetCollectResult 写的),
// title 和 shop_name 也要清:它们是**采集回来的**(SetCollectResult 写的),
// 属于采集结果的一部分。留着的话状态显示"未采集"、标题却有值,
// 操作员会以为已经采过了。
_, err := q.Exec(`
UPDATE pdd_products
SET deleted_at = NULL, url = ?, collect_status = 'pending',
title = NULL, skus_json = NULL, collect_msg = NULL,
title = NULL, shop_name = NULL, skus_json = NULL, collect_msg = NULL,
artifact_ref = NULL, collected_at = NULL, updated_at = ?
WHERE goods_id = ?`,
url, now, goodsID)
@@ -304,18 +305,21 @@ func SoftDeletePddProducts(q Execer, goodsIDs []string) (int64, error) {
//
// `[必须]` 按 **PDD 的 goods_id** 定位,不是蝦皮的。采集的对象是 PDD 商品。
//
// title 由调用方从采集结果里取出来传进来,用于人工核对"采的是不是要的那个商品"。
func SetCollectResult(q Execer, pddGoodsID, title, skusJSON string) error {
// title 和 shopName 由调用方从采集结果里取出来,用于人工核对商品来源。
// shopName 为空时保留数据库已有值:这次没采到不代表上次采到的店铺名失效。
func SetCollectResult(q Execer, pddGoodsID, title, shopName, skusJSON string) error {
if pddGoodsID == "" {
return fmt.Errorf("pdd goods_id 不能为空")
}
now := model.NowISO()
res, err := q.Exec(`
UPDATE pdd_products
SET skus_json = ?, title = ?, collect_status = 'collected',
SET skus_json = ?, title = ?,
shop_name = CASE WHEN TRIM(?) = '' THEN shop_name ELSE ? END,
collect_status = 'collected',
collect_msg = NULL, collected_at = ?, updated_at = ?
WHERE goods_id = ? AND deleted_at IS NULL`,
skusJSON, title, now, now, pddGoodsID)
skusJSON, title, shopName, shopName, now, now, pddGoodsID)
if err != nil {
return fmt.Errorf("保存 PDD 商品 %s 的采集结果失败: %w", pddGoodsID, err)
}