Files
cmautobuy/admin/repository/task.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

244 lines
7.8 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
import (
"database/sql"
"fmt"
"strings"
"cmautobuy/admin/model"
)
// claimCandidateLimit 是一次最多尝试抢多少条。
// 抢不到说明被别的客户端拿走了,再试下一条;都抢不到就当作没任务。
const claimCandidateLimit = 10
// ClaimNextTask 为指定客户端领取一个任务。
//
// 没有可领的任务时返回 (nil, nil) —— 调用方据此返回 204。
//
// 防并发的做法是**条件更新 + 检查影响行数**:先查出候选,
// 再用 `WHERE task_id = ? AND status = 'assigned'` 去更新,
// 影响行数为 0 就说明被别人抢先了,换下一条。
// 不用 SELECT ... FOR UPDATE,SQLite 没有那个。
func ClaimNextTask(db *sql.DB, clientID string, supportedTypes []string) (*model.Task, error) {
if clientID == "" {
return nil, fmt.Errorf("client_id 不能为空")
}
query := `SELECT task_id FROM tasks
WHERE assigned_client = ? AND status = 'assigned'`
args := []any{clientID}
// 客户端只声明支持某些类型时,不要给它别的类型
if len(supportedTypes) > 0 {
placeholders := strings.TrimSuffix(strings.Repeat("?,", len(supportedTypes)), ",")
query += ` AND task_type IN (` + placeholders + `)`
for _, t := range supportedTypes {
args = append(args, t)
}
}
query += ` ORDER BY priority DESC, created_at LIMIT ?`
args = append(args, claimCandidateLimit)
rows, err := db.Query(query, args...)
if err != nil {
return nil, fmt.Errorf("查询可领任务失败: %w", err)
}
var candidates []string
for rows.Next() {
var id string
if err := rows.Scan(&id); err != nil {
rows.Close()
return nil, fmt.Errorf("读取候选任务失败: %w", err)
}
candidates = append(candidates, id)
}
rows.Close()
if err := rows.Err(); err != nil {
return nil, err
}
now := model.NowISO()
for _, taskID := range candidates {
res, err := db.Exec(`
UPDATE tasks SET status = 'claimed', claimed_at = ?, updated_at = ?
WHERE task_id = ? AND status = 'assigned'`,
now, now, taskID)
if err != nil {
return nil, fmt.Errorf("领取任务 %s 失败: %w", taskID, err)
}
n, err := res.RowsAffected()
if err != nil {
return nil, err
}
if n == 0 {
continue // 被别的客户端抢先了,换下一条
}
// 记一笔领取历史。提交结果时要靠它判断这台客户端有没有领过——
// 任务重派后 assigned_client 会变,只看它就查不出来了。
if err := RecordClaim(db, taskID, clientID, now); err != nil {
return nil, err
}
return GetTask(db, taskID)
}
return nil, nil // 没有可领的任务
}
// GetTask 按编号读一条任务。
func GetTask(db *sql.DB, taskID string) (*model.Task, error) {
var t model.Task
var assigned, claimedAt, sybID, orderNo, goodsID, skuID sql.NullString
var pddGoodsID, pddOptions, resultData, errCode, errMsg, finishedAt sql.NullString
var quantity, maxPrice sql.NullInt64
err := db.QueryRow(`
SELECT task_id, task_type, status, version, priority,
assigned_client, claimed_at,
syb_id, order_no, goods_id, shopee_sku_id,
pdd_goods_url, pdd_goods_id, pdd_options,
quantity, max_price_cent,
result_data, error_code, error_message, finished_at,
created_at, updated_at
FROM tasks WHERE task_id = ?`, taskID).Scan(
&t.TaskID, &t.TaskType, &t.Status, &t.Version, &t.Priority,
&assigned, &claimedAt,
&sybID, &orderNo, &goodsID, &skuID,
&t.PddGoodsURL, &pddGoodsID, &pddOptions,
&quantity, &maxPrice,
&resultData, &errCode, &errMsg, &finishedAt,
&t.CreatedAt, &t.UpdatedAt)
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, fmt.Errorf("读取任务 %s 失败: %w", taskID, err)
}
t.AssignedClient = assigned.String
t.ClaimedAt = claimedAt.String
t.SybID = sybID.String
t.OrderNo = orderNo.String
t.GoodsID = goodsID.String
t.ShopeeSKUID = skuID.String
t.PddGoodsID = pddGoodsID.String
t.PddOptions = pddOptions.String
t.Quantity = int(quantity.Int64)
t.MaxPriceCent = maxPrice.Int64
t.ResultData = resultData.String
t.ErrorCode = errCode.String
t.ErrorMessage = errMsg.String
t.FinishedAt = finishedAt.String
return &t, nil
}
// RecordClaim 记一笔"某客户端领过某任务"。
//
// 同一台客户端重复领同一个任务时只更新时间,不报错。
func RecordClaim(q Execer, taskID, clientID, claimedAt string) error {
_, err := q.Exec(`
INSERT INTO task_claims (task_id, client_id, claimed_at)
VALUES (?, ?, ?)
ON CONFLICT(task_id, client_id) DO UPDATE SET claimed_at = excluded.claimed_at`,
taskID, clientID, claimedAt)
if err != nil {
return fmt.Errorf("记录领取历史失败 task=%s client=%s: %w", taskID, clientID, err)
}
return nil
}
// HasEverClaimed 判断这台客户端**曾经**领过这个任务。
//
// 用于提交结果时的权限判断:只有从没领过的才拒绝(403)。
// 任务已取消、已重派给别人,都**不影响**这个判断——
// 契约要求那些情况仍然要接受结果,见 docs/admin/04-client-api.md §4.1。
func HasEverClaimed(q Execer, taskID, clientID string) (bool, error) {
var one int
err := q.QueryRow(
`SELECT 1 FROM task_claims WHERE task_id = ? AND client_id = ?`,
taskID, clientID).Scan(&one)
if err == sql.ErrNoRows {
return false, nil
}
if err != nil {
return false, fmt.Errorf("查询领取历史失败: %w", err)
}
return true, nil
}
// TaskInfo 是提交结果时需要知道的任务基本信息。
type TaskInfo struct {
TaskType model.TaskType
// GoodsID 是关联的**蝦皮**商品编号。
GoodsID string
// PddGoodsID 是要采集/购买的**拼多多**商品编号。
// 采集结果落到 pdd_products 时用的是它,不是 GoodsID —— 被采集的是 PDD 商品。
PddGoodsID string
}
// GetTaskInfo 查任务的类型和关联商品。任务不存在时返回 (nil, nil)。
func GetTaskInfo(q Execer, taskID string) (*TaskInfo, error) {
var info TaskInfo
var gid, pddGID sql.NullString
err := q.QueryRow(
`SELECT task_type, goods_id, pdd_goods_id FROM tasks WHERE task_id = ?`,
taskID).Scan(&info.TaskType, &gid, &pddGID)
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, fmt.Errorf("查询任务 %s 失败: %w", taskID, err)
}
info.GoodsID = gid.String
info.PddGoodsID = pddGID.String
return &info, nil
}
// MarkTaskSucceeded 记录成功结果。
//
// `[必须]` **不检查任务当前状态**。任务已取消、已重派给别的客户端,
// 都照样接受——客户端可能真的已经下单了,这些数据必须能留痕。
// 理由见 docs/admin/04-client-api.md §4.1。
func MarkTaskSucceeded(q Execer, taskID, resultData string) error {
now := model.NowISO()
_, err := q.Exec(`
UPDATE tasks
SET status = 'succeeded', result_data = ?, finished_at = ?, updated_at = ?
WHERE task_id = ?`,
resultData, now, now, taskID)
if err != nil {
return fmt.Errorf("标记任务 %s 成功失败: %w", taskID, err)
}
return nil
}
// MarkTaskFailure 记录失败/需人工处理的结果。
//
// status 由客户端报告,映射规则见 docs/admin/04-client-api.md §5:
//
// retry_wait -> assigned(放回去等它再来领)
// manual_review -> manual_review
// failed -> failed
// cancelled -> cancelled
//
// 同样**不检查任务当前状态**,理由同 MarkTaskSucceeded。
func MarkTaskFailure(q Execer, taskID string, newStatus model.TaskStatus, errCode, errMsg string) error {
now := model.NowISO()
// 放回待领取的话不算结束,finished_at 保持为空
finishedAt := any(now)
if newStatus == model.TaskAssigned {
finishedAt = nil
}
_, err := q.Exec(`
UPDATE tasks
SET status = ?, error_code = ?, error_message = ?,
finished_at = ?, updated_at = ?
WHERE task_id = ?`,
newStatus, errCode, errMsg, finishedAt, now, taskID)
if err != nil {
return fmt.Errorf("标记任务 %s 失败状态出错: %w", taskID, err)
}
return nil
}