fix: 修复顺运宝采集任务孤儿状态 (#165)

This commit is contained in:
chengma
2026-08-11 17:00:33 +08:00
parent c8b1fd686a
commit 9e328984e9
17 changed files with 508 additions and 67 deletions
+104 -2
View File
@@ -19,7 +19,7 @@ import (
"cmautobuy/admin/spec"
)
const mysqlSchemaVersion = 11
const mysqlSchemaVersion = 12
// OpenMySQL 打开生产 MySQL 8 数据库。错误信息绝不包含完整 DSN 或密码。
func OpenMySQL(cfg config.DatabaseConfig) (*sql.DB, error) {
@@ -555,10 +555,102 @@ func MigrateMySQL(db *sql.DB) error {
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 11, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
return fmt.Errorf("记录 MySQL schema v11 失败: %w", err)
}
current = 11
}
if current < 12 {
if err := migrateMySQLV12(db); err != nil {
return fmt.Errorf("执行 MySQL schema v12 失败: %w", err)
}
if err := checkMySQLV12Shape(db); err != nil {
return fmt.Errorf("MySQL schema v12 自检失败,未记录版本: %w", err)
}
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 12, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
return fmt.Errorf("记录 MySQL schema v12 失败: %w", err)
}
}
return CheckMySQLSchema(db)
}
// migrateMySQLV12 保存采集任务的顺运宝来源,并回收没有有效任务的孤儿采集中状态。
// CREATE TABLE IF NOT EXISTS 和带条件的 UPDATE 都可重放,适合 MySQL DDL 隐式提交后的重启恢复。
func migrateMySQLV12(db *sql.DB) error {
if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS task_syb_sources (
task_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
syb_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
created_at VARCHAR(35) NOT NULL,
PRIMARY KEY (task_id, syb_id),
KEY idx_task_syb_sources_syb (syb_id, task_id),
CONSTRAINT fk_task_syb_sources_task FOREIGN KEY (task_id)
REFERENCES tasks(task_id) ON DELETE CASCADE,
CONSTRAINT fk_task_syb_sources_syb FOREIGN KEY (syb_id)
REFERENCES syb_orders(syb_id) ON DELETE CASCADE
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`); err != nil {
return fmt.Errorf("建立采集任务顺运宝来源表失败: %w", err)
}
now := model.NowISO()
if _, err := db.Exec(`UPDATE pdd_products AS pp
SET collect_status='pending', collect_msg=NULL, updated_at=?
WHERE pp.collect_status='collecting' AND pp.deleted_at IS NULL
AND NOT EXISTS (
SELECT 1 FROM tasks t
WHERE t.task_type='collect' AND t.pdd_goods_id=pp.goods_id
AND t.status IN ('pending','assigned','claimed')
)`, now); err != nil {
return fmt.Errorf("回收无有效任务的 PDD 采集中状态失败: %w", err)
}
return nil
}
func checkMySQLV12Shape(db *sql.DB) error {
exists, err := mysqlTableExists(db, "task_syb_sources")
if err != nil || !exists {
return fmt.Errorf("采集任务顺运宝来源表缺失")
}
for _, column := range []struct {
name string
length int64
collation string
}{
{"task_id", 191, "utf8mb4_bin"},
{"syb_id", 191, "utf8mb4_bin"},
{"created_at", 35, "utf8mb4_0900_ai_ci"},
} {
if err := checkMySQLVarcharColumn(db, "task_syb_sources", column.name, column.length, false, column.collation, ""); err != nil {
return err
}
}
for _, index := range []struct{ name, columns string }{
{"PRIMARY", "task_id,syb_id"},
{"idx_task_syb_sources_syb", "syb_id,task_id"},
} {
var columns string
if err := db.QueryRow(`SELECT GROUP_CONCAT(column_name ORDER BY seq_in_index)
FROM information_schema.statistics
WHERE table_schema=DATABASE() AND table_name='task_syb_sources' AND index_name=?`, index.name).Scan(&columns); err != nil || columns != index.columns {
return fmt.Errorf("采集任务顺运宝来源索引 %s 不正确", index.name)
}
}
for _, foreignKey := range []struct{ name, column, table, target string }{
{"fk_task_syb_sources_task", "task_id", "tasks", "task_id"},
{"fk_task_syb_sources_syb", "syb_id", "syb_orders", "syb_id"},
} {
var column, table, target string
err := db.QueryRow(`SELECT column_name,referenced_table_name,referenced_column_name
FROM information_schema.key_column_usage
WHERE constraint_schema=DATABASE() AND table_name='task_syb_sources' AND constraint_name=?`, foreignKey.name).
Scan(&column, &table, &target)
if err != nil || column != foreignKey.column || table != foreignKey.table || target != foreignKey.target {
return fmt.Errorf("采集任务顺运宝来源约束 %s 不正确", foreignKey.name)
}
var deleteRule string
if err := db.QueryRow(`SELECT delete_rule FROM information_schema.referential_constraints
WHERE constraint_schema=DATABASE() AND table_name='task_syb_sources' AND constraint_name=?`, foreignKey.name).Scan(&deleteRule); err != nil || deleteRule != "CASCADE" {
return fmt.Errorf("采集任务顺运宝来源约束 %s 删除规则不正确", foreignKey.name)
}
}
return nil
}
// migrateMySQLV11 增加顺运宝店铺列,并从保留的原始 JSON 回填历史数据。
// DDL 和 UPDATE 都可重放:中断后再次启动只补缺列和仍为空的记录。
func migrateMySQLV11(db *sql.DB) error {
@@ -1167,6 +1259,7 @@ func CheckMySQLSchema(db *sql.DB) error {
"users", "web_sessions", "client_user_assignments", "syb_sync_runs", "admin_initialization_lock",
"spec_mapping_decisions",
"catalog_import_runs",
"task_syb_sources",
}
if err := checkMySQLSchema(db, mysqlRequiredTables); err != nil {
return err
@@ -1189,7 +1282,16 @@ func CheckMySQLSchema(db *sql.DB) error {
if err := checkMySQLV8Shape(db); err != nil {
return err
}
return checkMySQLV9Shape(db)
if err := checkMySQLV9Shape(db); err != nil {
return err
}
if err := checkMySQLV10Shape(db); err != nil {
return err
}
if err := checkMySQLV11Shape(db); err != nil {
return err
}
return checkMySQLV12Shape(db)
}
func checkMySQLV9Shape(db *sql.DB) error {
@@ -585,6 +585,71 @@ func TestMySQLMigrate_V11回填顺运宝店铺且可重放(t *testing.T) {
}
}
func TestMySQLMigrate_V11升级V12并修复孤儿采集中状态(t *testing.T) {
db := openMySQLMigrationTestDB(t)
defer db.Close()
cleanMySQLTestSchema(t, db)
defer cleanMySQLTestSchema(t, db)
if err := MigrateMySQL(db); err != nil {
t.Fatal(err)
}
// 模拟已在 v11 的生产库:去掉 v12 版本及其附加表,保留全部 v1-v11 结构。
mustExec(t, db, `DELETE FROM schema_migrations WHERE version=12`)
mustExec(t, db, `DROP TABLE task_syb_sources`)
now := model.NowISO()
mustExec(t, db, `INSERT INTO pdd_products(goods_id,url,collect_status,created_at,updated_at) VALUES
('ORPHAN','https://mobile.yangkeduo.com/goods.html?goods_id=1','collecting',?,?),
('ACTIVE','https://mobile.yangkeduo.com/goods.html?goods_id=2','collecting',?,?)`, now, now, now, now)
mustExec(t, db, `INSERT INTO tasks(task_id,task_type,status,pdd_goods_url,pdd_goods_id,created_at,updated_at)
VALUES('COL-ACTIVE','collect','pending','https://mobile.yangkeduo.com/goods.html?goods_id=2','ACTIVE',?,?)`, now, now)
if err := MigrateMySQL(db); err != nil {
t.Fatalf("v11 升级 v12 失败: %v", err)
}
if err := MigrateMySQL(db); err != nil {
t.Fatalf("v12 重放失败: %v", err)
}
var orphanStatus, activeStatus string
if err := db.QueryRow(`SELECT collect_status FROM pdd_products WHERE goods_id='ORPHAN'`).Scan(&orphanStatus); err != nil {
t.Fatal(err)
}
if err := db.QueryRow(`SELECT collect_status FROM pdd_products WHERE goods_id='ACTIVE'`).Scan(&activeStatus); err != nil {
t.Fatal(err)
}
if orphanStatus != "pending" || activeStatus != "collecting" {
t.Fatalf("v12 状态修复错误:orphan=%s active=%s", orphanStatus, activeStatus)
}
var versionCount int
if err := db.QueryRow(`SELECT COUNT(*) FROM schema_migrations WHERE version=12`).Scan(&versionCount); err != nil || versionCount != 1 {
t.Fatalf("v12 版本记录错误: count=%d err=%v", versionCount, err)
}
}
func TestMySQLMigrate_V12形状错误不记录版本(t *testing.T) {
db := openMySQLMigrationTestDB(t)
defer db.Close()
cleanMySQLTestSchema(t, db)
defer cleanMySQLTestSchema(t, db)
if err := MigrateMySQL(db); err != nil {
t.Fatal(err)
}
mustExec(t, db, `DELETE FROM schema_migrations WHERE version=12`)
mustExec(t, db, `DROP TABLE task_syb_sources`)
mustExec(t, db, `CREATE TABLE task_syb_sources (
task_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
syb_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
created_at VARCHAR(35) NOT NULL,
PRIMARY KEY(syb_id,task_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`)
if err := MigrateMySQL(db); err == nil {
t.Fatal("错误主键和缺失外键必须阻止 v12")
}
var versionCount int
if err := db.QueryRow(`SELECT COUNT(*) FROM schema_migrations WHERE version=12`).Scan(&versionCount); err != nil || versionCount != 0 {
t.Fatalf("v12 自检失败时不得记录版本:count=%d err=%v", versionCount, err)
}
}
func openMySQLMigrationTestDB(t *testing.T) *sql.DB {
t.Helper()
if os.Getenv("CMAUTOBUY_MYSQL_TEST") != "1" {
+20 -7
View File
@@ -387,7 +387,8 @@ func SetCollectFailed(q Execer, pddGoodsID, msg, artifactRef string) error {
//
// 允许发起采集的状态:
// - pending / failed —— 没有任务在跑;
// - collecting **且已经超时**(`updated_at` 早于 `now - model.CollectStaleAfter`)
// - collecting **且已经超时**(`updated_at` 早于 `now - model.CollectStaleAfter`),
// 或数据库里已经没有有效采集任务;
// —— 客户端离线、崩溃或任务被删都会让一个 collecting 卡住不动,
// 这些是常态不是异常,界面上必须有出口,见 #24。
//
@@ -414,11 +415,17 @@ func MarkCollecting(q Execer, pddGoodsID string) (bool, error) {
staleBefore := now.Add(-model.CollectStaleAfter).UTC().Format(model.TimeLayout)
res, err := q.Exec(`
UPDATE pdd_products
UPDATE pdd_products AS pp
SET collect_status = 'collecting', updated_at = ?
WHERE goods_id = ? AND deleted_at IS NULL
WHERE pp.goods_id = ? AND deleted_at IS NULL
AND ( collect_status IN ('pending', 'failed')
OR (collect_status = 'collecting' AND updated_at < ?) )`,
OR (collect_status = 'collecting' AND (
updated_at < ? OR NOT EXISTS (
SELECT 1 FROM tasks t
WHERE t.task_type='collect' AND t.pdd_goods_id=pp.goods_id
AND t.status IN ('pending','assigned','claimed')
)
)) )`,
nowISO, pddGoodsID, staleBefore)
if err != nil {
return false, fmt.Errorf("标记 PDD 商品 %s 采集中失败: %w", pddGoodsID, err)
@@ -442,11 +449,17 @@ func MarkRecollecting(q Execer, pddGoodsID string) (bool, error) {
staleBefore := now.Add(-model.CollectStaleAfter).UTC().Format(model.TimeLayout)
res, err := q.Exec(`
UPDATE pdd_products
UPDATE pdd_products AS pp
SET collect_status = 'collecting', updated_at = ?
WHERE goods_id = ? AND deleted_at IS NULL
WHERE pp.goods_id = ? AND deleted_at IS NULL
AND ( collect_status IN ('pending', 'failed', 'collected')
OR (collect_status = 'collecting' AND updated_at < ?) )`,
OR (collect_status = 'collecting' AND (
updated_at < ? OR NOT EXISTS (
SELECT 1 FROM tasks t
WHERE t.task_type='collect' AND t.pdd_goods_id=pp.goods_id
AND t.status IN ('pending','assigned','claimed')
)
)) )`,
nowISO, pddGoodsID, staleBefore)
if err != nil {
return false, fmt.Errorf("标记 PDD 商品 %s 重新采集中失败: %w", pddGoodsID, err)
+21 -13
View File
@@ -396,17 +396,18 @@ const sybOrderContextFrom = `
// SybOrderContext 是顺运宝明细及其当前蝦皮/PDD 处理上下文。
// 处理阶段由 service 计算,Repository 只提供数据库事实。
type SybOrderContext struct {
Order model.SybOrder
ShopeeExists bool
PddGoodsID string
PddGoodsURL string
PddCollectStatus string
PddCollectMsg string
PddSkusJSON string
PddUpdatedAt string
MappingOptionKey string
MappingOptions string
HasActiveTask bool
Order model.SybOrder
ShopeeExists bool
PddGoodsID string
PddGoodsURL string
PddCollectStatus string
PddCollectMsg string
PddSkusJSON string
PddUpdatedAt string
MappingOptionKey string
MappingOptions string
HasActiveTask bool
HasActiveCollectTask bool
}
func scanSybOrderContext(s rowScanner) (SybOrderContext, error) {
@@ -416,13 +417,13 @@ func scanSybOrderContext(s rowScanner) (SybOrderContext, error) {
var shopeeExists int
var pddGoodsID, pddGoodsURL, collectStatus, collectMsg, skusJSON, pddUpdatedAt sql.NullString
var mappingKey, mappingOptions sql.NullString
var hasActiveTask int
var hasActiveTask, hasActiveCollectTask int
err := s.Scan(
&c.Order.SybID, &c.Order.OrderNo, &shopName, &title, &productSpec, &specKey, &shopeeGoodsID,
&c.Order.Quantity, &priceCent, &imageURL, &c.Order.SybData,
&c.Order.CreatedAt, &c.Order.UpdatedAt, &shopeeExists,
&pddGoodsID, &pddGoodsURL, &collectStatus, &collectMsg, &skusJSON, &pddUpdatedAt,
&mappingKey, &mappingOptions, &hasActiveTask,
&mappingKey, &mappingOptions, &hasActiveTask, &hasActiveCollectTask,
)
c.Order.ShopName = shopName.String
c.Order.Title = title.String
@@ -441,6 +442,7 @@ func scanSybOrderContext(s rowScanner) (SybOrderContext, error) {
c.MappingOptionKey = mappingKey.String
c.MappingOptions = mappingOptions.String
c.HasActiveTask = hasActiveTask != 0
c.HasActiveCollectTask = hasActiveCollectTask != 0
return c, err
}
@@ -456,6 +458,9 @@ func ListSybOrderContexts(q Execer, filter SybOrderFilter, limit, offset int) ([
pp.skus_json, pp.updated_at, sm.pdd_option_key, sm.pdd_options,
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
AND t.syb_id = so.syb_id
AND t.status IN ('pending', 'assigned', 'claimed')),
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'collect'
AND t.pdd_goods_id = pp.goods_id
AND t.status IN ('pending', 'assigned', 'claimed'))` +
sybOrderContextFrom + where + `
ORDER BY so.updated_at DESC, so.syb_id DESC`
@@ -490,6 +495,9 @@ func GetSybOrderContext(q Execer, sybID string) (*SybOrderContext, error) {
pp.skus_json, pp.updated_at, sm.pdd_option_key, sm.pdd_options,
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
AND t.syb_id = so.syb_id
AND t.status IN ('pending', 'assigned', 'claimed')),
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'collect'
AND t.pdd_goods_id = pp.goods_id
AND t.status IN ('pending', 'assigned', 'claimed'))`+
sybOrderContextFrom+` WHERE so.syb_id = ?`, sybID))
if errors.Is(err, sql.ErrNoRows) {
+86 -2
View File
@@ -379,8 +379,14 @@ func taskFilterClause(filter TaskFilter) (string, []any) {
if kw := strings.TrimSpace(filter.Keyword); kw != "" {
pattern := "%" + escapeLike(kw) + "%"
clauses = append(clauses,
"(t.task_id LIKE ? ESCAPE '!' OR t.order_no LIKE ? ESCAPE '!' OR t.pdd_goods_id LIKE ? ESCAPE '!')")
args = append(args, pattern, pattern, pattern)
`(t.task_id LIKE ? ESCAPE '!' OR t.order_no LIKE ? ESCAPE '!' OR t.pdd_goods_id LIKE ? ESCAPE '!'
OR EXISTS (
SELECT 1 FROM task_syb_sources tss
JOIN syb_orders so ON so.syb_id=tss.syb_id
WHERE tss.task_id=t.task_id
AND (tss.syb_id LIKE ? ESCAPE '!' OR so.order_no LIKE ? ESCAPE '!')
))`)
args = append(args, pattern, pattern, pattern, pattern, pattern)
}
if filter.VisibleUserID != "" {
clauses = append(clauses, "t.created_by_user_id = ?")
@@ -529,6 +535,59 @@ func DeleteTasksInScope(q Execer, taskIDs []string, visibleUserID string) (int64
return res.RowsAffected()
}
// ListCollectTaskGoodsIDsInScope 在删除前找出受影响的 PDD 商品。
// 调用方必须和删除使用同一事务、同一可见范围,避免越权数据影响状态回收。
func ListCollectTaskGoodsIDsInScope(q Execer, taskIDs []string, visibleUserID string) ([]string, error) {
if len(taskIDs) == 0 {
return nil, nil
}
placeholders := strings.TrimSuffix(strings.Repeat("?,", len(taskIDs)), ",")
args := make([]any, 0, len(taskIDs)+1)
for _, id := range taskIDs {
args = append(args, id)
}
where := `task_id IN (` + placeholders + `) AND task_type='collect' AND pdd_goods_id IS NOT NULL`
if visibleUserID != "" {
where += ` AND created_by_user_id = ?`
args = append(args, visibleUserID)
}
rows, err := q.Query(`SELECT DISTINCT pdd_goods_id FROM tasks WHERE `+where, args...)
if err != nil {
return nil, fmt.Errorf("查询待删除采集任务的 PDD 商品失败: %w", err)
}
defer rows.Close()
var goodsIDs []string
for rows.Next() {
var goodsID string
if err := rows.Scan(&goodsID); err != nil {
return nil, fmt.Errorf("读取待删除采集任务的 PDD 商品失败: %w", err)
}
goodsIDs = append(goodsIDs, goodsID)
}
return goodsIDs, rows.Err()
}
// ResetCollectingIfNoActiveTask 把没有有效采集任务的孤儿状态回收到待采集。
// 条件判断和更新在一条 SQL 内完成,避免并发建任务时误覆盖 collecting。
func ResetCollectingIfNoActiveTask(q Execer, goodsID string) (bool, error) {
result, err := q.Exec(`UPDATE pdd_products AS pp
SET collect_status='pending', collect_msg=NULL, updated_at=?
WHERE pp.goods_id=? AND pp.collect_status='collecting'
AND NOT EXISTS (
SELECT 1 FROM tasks t
WHERE t.task_type='collect' AND t.pdd_goods_id=pp.goods_id
AND t.status IN ('pending','assigned','claimed')
)`, model.NowISO(), goodsID)
if err != nil {
return false, fmt.Errorf("回收商品 %s 的孤立采集中状态失败: %w", goodsID, err)
}
n, err := result.RowsAffected()
if err != nil {
return false, fmt.Errorf("确认商品 %s 的采集状态回收结果失败: %w", goodsID, err)
}
return n == 1, nil
}
// InsertCollectTask 建一条不指定客户端的采集任务。
// 保留这个入口给蝦皮、顺运宝等现有流程使用,避免它们被 PDD 页的新选项影响。
func InsertCollectTask(q Execer, taskID, goodsID, goodsURL string) error {
@@ -567,6 +626,31 @@ func InsertCollectTaskForClientAndUser(q Execer, taskID, goodsID, goodsURL, assi
return nil
}
// InsertCollectTaskSybSources 保存一条采集任务对应的全部顺运宝明细来源。
// 同一个 PDD 商品可能由多条明细共同发起,来源不能压成 tasks 上的单列。
func InsertCollectTaskSybSources(q Execer, taskID string, sybIDs []string) error {
taskID = strings.TrimSpace(taskID)
if taskID == "" {
return fmt.Errorf("采集任务编号不能为空")
}
seen := make(map[string]struct{}, len(sybIDs))
for _, raw := range sybIDs {
sybID := strings.TrimSpace(raw)
if sybID == "" {
return fmt.Errorf("顺运宝明细编号不能为空")
}
if _, ok := seen[sybID]; ok {
continue
}
seen[sybID] = struct{}{}
if _, err := q.Exec(`INSERT INTO task_syb_sources(task_id,syb_id,created_at)
VALUES(?,?,?) ON DUPLICATE KEY UPDATE created_at=created_at`, taskID, sybID, model.NowISO()); err != nil {
return fmt.Errorf("保存采集任务 %s 的顺运宝来源 %s 失败: %w", taskID, sybID, err)
}
}
return nil
}
// HasActivePurchaseTask 判断顺运宝明细是否已有尚未结束的采购任务。
func HasActivePurchaseTask(q Execer, sybID string) (bool, error) {
var count int