Files
cmautobuy/admin/repository/task.go
T
chengmaandClaude Opus 5 dd387e3ee5 feat: claim 支持领取无主任务 (#17)
PDD 商品页要能「创建采集任务」而不指定客户端——采集只是浏览商品页,
没有副作用,哪台设备采都一样。但原来的查询是

    WHERE assigned_client = ? AND status = 'assigned'

无主任务(assigned_client 为空 + pending)永远没人能领,建出来就是死的。

改动
- ClaimNextTask 同时查两种:指定给本机的 + 无主的
- 排序 ORDER BY (assigned_client IS NULL), priority DESC, created_at
  ——指定给本机的优先。显式分配是人为决定,应当先兑现
- 原子更新两种情况合成一条语句:对"指定给我的"写 assigned_client
  是写同一个值无副作用;对无主的,这一步就是"谁领到就标记谁"
- 表结构不用动(assigned_client 本来可空,status 已有 pending)

推翻了一条已定案的规则
Client 契约 §5.1 原写「Admin 只把任务分配给指定的 Client」,
现改为两种并存并说明各自适用场景:
- 采集任务不指定客户端
- 采购任务可指定可留空。涉及钱和账号——不同设备可能登着不同的
  拼多多账号,需要指定账号时必须显式分配,留空即接受"谁先抢到谁下单"

Client 侧对两种没有区别,不需要知道任务原来有没有主。

已验证(Go 1.23.0)
- 新增 7 个测试,全量 62 个全过
- 并发抢占用例重复 20 次稳定:8 个客户端抢同一条无主任务,
  正好 1 个拿到,且 assigned_client 记的就是那个赢家
- 既有测试未受影响,"只领分配给自己的"仍然成立

一处仍未解决的风险(已记入 #17 风险表)
tasks 表没有字段标记"该任务需要真实下单",所以契约里
"不向 dry_run 客户端分配真实下单任务"实际无法执行。
真实下单开关关闭时不出问题,开启前必须补该字段。

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

265 lines
8.9 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。
//
// # 两种任务都能领
//
// 指定给本机的 assigned_client = 我 且 status = 'assigned'
// 无主的 assigned_client 为空 且 status = 'pending'
//
// 采集任务创建时不指定客户端(浏览商品页没有副作用,哪台设备采都一样),
// 所以必须支持第二种,否则那些任务永远没人能领。
// 采购任务可以指定也可以留空——涉及钱和账号时应当显式分配。
//
// **指定给本机的优先。** 显式分配是人为决定,应当先兑现;
// 无主任务谁抢都一样,可以等。
//
// # 防并发
//
// 做法是**条件更新 + 检查影响行数**:先查出候选,再带着状态条件去更新,
// 影响行数为 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')
OR (assigned_client IS NULL AND status = 'pending') )`
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)
}
}
// (assigned_client IS NULL) 为 0/1,0 排前面 —— 指定给本机的优先于无主的
query += ` ORDER BY (assigned_client IS NULL), 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 {
// 两种情况合成一条语句:对"指定给我的"那种,写 assigned_client
// 是写同一个值,无副作用;对无主的,这一步就是"谁领到就标记谁"。
res, err := db.Exec(`
UPDATE tasks
SET status = 'claimed', assigned_client = ?,
claimed_at = ?, updated_at = ?
WHERE task_id = ?
AND ( (status = 'assigned' AND assigned_client = ?)
OR (status = 'pending' AND assigned_client IS NULL) )`,
clientID, now, now, taskID, clientID)
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
}