diff --git a/admin/repository/task.go b/admin/repository/task.go index 3663cd4..6cb7b87 100644 --- a/admin/repository/task.go +++ b/admin/repository/task.go @@ -16,8 +16,21 @@ const claimCandidateLimit = 10 // // 没有可领的任务时返回 (nil, nil) —— 调用方据此返回 204。 // -// 防并发的做法是**条件更新 + 检查影响行数**:先查出候选, -// 再用 `WHERE task_id = ? AND status = 'assigned'` 去更新, +// # 两种任务都能领 +// +// 指定给本机的 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) { @@ -26,7 +39,8 @@ func ClaimNextTask(db *sql.DB, clientID string, supportedTypes []string) (*model } query := `SELECT task_id FROM tasks - WHERE assigned_client = ? AND status = 'assigned'` + WHERE ( (assigned_client = ? AND status = 'assigned') + OR (assigned_client IS NULL AND status = 'pending') )` args := []any{clientID} // 客户端只声明支持某些类型时,不要给它别的类型 @@ -37,7 +51,8 @@ func ClaimNextTask(db *sql.DB, clientID string, supportedTypes []string) (*model args = append(args, t) } } - query += ` ORDER BY priority DESC, created_at LIMIT ?` + // (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...) @@ -60,10 +75,16 @@ func ClaimNextTask(db *sql.DB, clientID string, supportedTypes []string) (*model now := model.NowISO() for _, taskID := range candidates { + // 两种情况合成一条语句:对"指定给我的"那种,写 assigned_client + // 是写同一个值,无副作用;对无主的,这一步就是"谁领到就标记谁"。 res, err := db.Exec(` - UPDATE tasks SET status = 'claimed', claimed_at = ?, updated_at = ? - WHERE task_id = ? AND status = 'assigned'`, - now, now, taskID) + 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) } diff --git a/admin/service/claim_unassigned_test.go b/admin/service/claim_unassigned_test.go new file mode 100644 index 0000000..661d7e2 --- /dev/null +++ b/admin/service/claim_unassigned_test.go @@ -0,0 +1,198 @@ +package service + +import ( + "database/sql" + "sync" + "testing" + + "cmautobuy/admin/model" +) + +// insertUnassignedTask 插一条**无主**任务:没有指定客户端,状态是 pending。 +// 采集任务就长这样——创建时不指定客户端,谁领到算谁的。 +func insertUnassignedTask(t *testing.T, db *sql.DB, taskID, taskType string) { + t.Helper() + now := model.NowISO() + _, err := db.Exec(` + INSERT INTO tasks (task_id, task_type, status, assigned_client, + pdd_goods_url, pdd_goods_id, created_at, updated_at) + VALUES (?, ?, 'pending', NULL, 'https://x/1', '1', ?, ?)`, + taskID, taskType, now, now) + if err != nil { + t.Fatalf("插入无主任务失败: %v", err) + } +} + +func taskAssignee(t *testing.T, db *sql.DB, taskID string) (client, status string) { + t.Helper() + var c sql.NullString + if err := db.QueryRow( + `SELECT assigned_client, status FROM tasks WHERE task_id = ?`, taskID, + ).Scan(&c, &status); err != nil { + t.Fatalf("查询任务失败: %v", err) + } + return c.String, status +} + +// ── 无主任务能被领到 ─────────────────────────────────── + +func TestClaim_无主任务能被领到并标记领取者(t *testing.T) { + db := newTestDB(t) + insertUnassignedTask(t, db, "T-FREE", "collect") + + task, err := ClaimNextTask(db, "client-001", []string{"collect"}) + if err != nil { + t.Fatalf("领取失败: %v", err) + } + if task == nil { + t.Fatal("无主任务应该能被领到——不然采集任务建出来就是死的") + } + + client, status := taskAssignee(t, db, "T-FREE") + if client != "client-001" { + t.Errorf("领取后应把领取者写进 assigned_client,实际 %q", client) + } + if status != "claimed" { + t.Errorf("状态应为 claimed,实际 %s", status) + } + if task.ClaimedAt == "" { + t.Error("claimed_at 应该有值") + } +} + +func TestClaim_无主任务只能被领一次(t *testing.T) { + db := newTestDB(t) + insertUnassignedTask(t, db, "T-FREE", "collect") + + if first, _ := ClaimNextTask(db, "client-001", []string{"collect"}); first == nil { + t.Fatal("第一次应该领到") + } + second, err := ClaimNextTask(db, "client-002", []string{"collect"}) + if err != nil { + t.Fatalf("第二次领取报错: %v", err) + } + if second != nil { + t.Errorf("已被领走的任务不该再被领到:%s", second.TaskID) + } +} + +// ── 优先级:指定给我的优先于无主的 ───────────────────── + +func TestClaim_指定给本机的优先于无主的(t *testing.T) { + db := newTestDB(t) + // 先插无主的,让它 created_at 更早——如果没有优先级规则, + // 按 created_at 排序会先拿到它,测试就能发现问题 + insertUnassignedTask(t, db, "T-FREE", "purchase") + insertTask(t, db, "T-MINE", "client-001") + + task, err := ClaimNextTask(db, "client-001", []string{"purchase"}) + if err != nil { + t.Fatalf("领取失败: %v", err) + } + if task == nil { + t.Fatal("应该领到任务") + } + if task.TaskID != "T-MINE" { + t.Errorf("应优先领取指定给本机的 T-MINE,实际领到 %s —— "+ + "显式分配是人为决定,应当先兑现", task.TaskID) + } +} + +func TestClaim_指定给别人的仍然领不到(t *testing.T) { + db := newTestDB(t) + insertTask(t, db, "T-OTHER", "client-999") + + task, err := ClaimNextTask(db, "client-001", []string{"purchase"}) + if err != nil { + t.Fatalf("领取失败: %v", err) + } + if task != nil { + t.Errorf("不该领到指定给别的客户端的任务:%s", task.TaskID) + } +} + +// ── 类型过滤仍然生效 ─────────────────────────────────── + +func TestClaim_无主任务也受supported_types约束(t *testing.T) { + db := newTestDB(t) + insertUnassignedTask(t, db, "T-FREE", "purchase") + + task, err := ClaimNextTask(db, "client-001", []string{"collect"}) + if err != nil { + t.Fatalf("领取失败: %v", err) + } + if task != nil { + t.Errorf("只声明 collect 的客户端不该拿到 purchase 任务:%s", task.TaskID) + } +} + +// ── 并发抢占 ─────────────────────────────────────────── + +// 多个客户端同时抢同一条无主任务,只能有一个拿到。 +// 这是"谁领到算谁的"这个模式的安全底线——两台都拿到就会重复采集/重复下单。 +func TestClaim_并发抢无主任务只有一个拿到(t *testing.T) { + db := newTestDB(t) + insertUnassignedTask(t, db, "T-ONLY-ONE", "collect") + + const workers = 8 + var ( + wg sync.WaitGroup + mu sync.Mutex + winners []string + lastErr error + ) + for i := 0; i < workers; i++ { + clientID := "client-" + string(rune('A'+i)) + wg.Add(1) + go func() { + defer wg.Done() + task, err := ClaimNextTask(db, clientID, []string{"collect"}) + mu.Lock() + defer mu.Unlock() + if err != nil { + lastErr = err + return + } + if task != nil { + winners = append(winners, clientID) + } + }() + } + wg.Wait() + + if lastErr != nil { + t.Fatalf("并发领取出错: %v", lastErr) + } + if len(winners) != 1 { + t.Fatalf("同一条无主任务被 %d 个客户端领到,期望正好 1 个:%v", + len(winners), winners) + } + + // 而且库里记的领取者必须就是那个赢家 + client, _ := taskAssignee(t, db, "T-ONLY-ONE") + if client != winners[0] { + t.Errorf("assigned_client 记的是 %q,但实际领到的是 %q", client, winners[0]) + } +} + +// ── 领取历史仍然被记录 ───────────────────────────────── + +// 提交结果时的权限判断依赖 task_claims(只有从没领过的才 403), +// 无主任务这条路径也必须记。 +func TestClaim_无主任务领取后也记领取历史(t *testing.T) { + db := newTestDB(t) + insertUnassignedTask(t, db, "T-FREE", "collect") + RegisterClient(db, model.Client{ClientID: "client-001"}, true) + + if _, err := ClaimNextTask(db, "client-001", []string{"collect"}); err != nil { + t.Fatalf("领取失败: %v", err) + } + + var n int + db.QueryRow(`SELECT COUNT(*) FROM task_claims + WHERE task_id = 'T-FREE' AND client_id = 'client-001'`).Scan(&n) + if n != 1 { + t.Errorf("领取历史应有 1 条,实际 %d —— "+ + "没有它的话提交结果会被误判成 403", n) + } +} diff --git a/docs/admin/04-client-api.md b/docs/admin/04-client-api.md index 121a9d9..4d2e705 100644 --- a/docs/admin/04-client-api.md +++ b/docs/admin/04-client-api.md @@ -58,8 +58,12 @@ X-Client-Id: client-001 **处理步骤:** 1. 先做客户端注册/更新(§3); -2. 查 `tasks`:`assigned_client = X-Client-Id` 且 `status = 'assigned'`, - 按 `priority DESC, created_at` 取**一条**; +2. 查 `tasks`,**两种任务都要查**: + - 指定给本机的:`assigned_client = X-Client-Id` 且 `status = 'assigned'` + - 无主的:`assigned_client IS NULL` 且 `status = 'pending'` + + 排序 `ORDER BY (assigned_client IS NULL), priority DESC, created_at` + —— **指定给本机的优先**; 3. 没有就返回 `204 No Content`; 4. 有就把 `status` 改成 `claimed`、记 `claimed_at`,返回任务。 @@ -69,8 +73,12 @@ X-Client-Id: client-001 `[必须]` 第 4 步的查询和更新要在**同一个事务**里,用条件更新防并发: ```sql -UPDATE tasks SET status = 'claimed', claimed_at = ?, updated_at = ? -WHERE task_id = ? AND status = 'assigned'; +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) ); +-- 两种情况合成一条:对"指定给我的",写 assigned_client 是写同一个值,无副作用。 -- 检查影响行数,为 0 说明被别人抢先了,重新取下一条 ``` @@ -268,7 +276,10 @@ Client 契约里散落的 Admin 侧硬要求,汇总在这里,**可以直接 - [x] claim 隐式登记**不更新名称** - [x] 登记和 claim **共用同一套字段校验**,非法内容返回 `422 INVALID_CLIENT_PROFILE` - [x] 校验不通过时不写库 -- [ ] `claim` 一次只返回一个任务,且只返回分配给该客户端的 +- [ ] `claim` 一次只返回一个任务 +- [ ] 能领到无主任务(`assigned_client IS NULL` + `pending`),领取时写入领取者 +- [ ] 指定给本机的任务优先于无主任务 +- [ ] 不会返回指定给别的客户端的任务 - [ ] `claim` 原子完成,并发下不会把同一任务发给两个客户端 - [ ] `claim` 无任务时返回 `204`,不是 `200` 加空对象 - [ ] 响应不含租约、不含 Admin 侧状态 diff --git a/docs/client/04-admin-api-contract.md b/docs/client/04-admin-api-contract.md index 7d5cf82..e1155d1 100644 --- a/docs/client/04-admin-api-contract.md +++ b/docs/client/04-admin-api-contract.md @@ -216,7 +216,26 @@ POST /api/v1/client/tasks/claim ### 5.1 Admin 侧的分配语义 -已确认:**Admin 只把任务分配给指定的 Client**,`claim` 只会返回分配给本 `X-Client-Id` 的任务。 +任务分两种,`claim` **两种都会返回**: + +| 类型 | `assigned_client` | Admin 侧状态 | 谁能领 | +|---|---|---|---| +| **指定分配** | 某个 Client 编号 | `assigned` | 只有那个 Client | +| **无主** | 空 | `pending` | **谁先抢到算谁的**,领取时才记下领取者 | + +`[必须]` **指定给本机的优先于无主的。** 显式分配是人为决定,应当先兑现; +无主任务谁抢都一样,可以等。 + +哪种任务用哪种方式: + +- **采集任务不指定客户端。** 采集只是浏览商品页,没有副作用,哪台设备采都一样, + 没必要每次都挑一台。 +- **采购任务可以指定,也允许留空。** 涉及钱和账号——不同设备可能登着不同的 + 拼多多账号,买到谁头上是有区别的。**需要指定账号时必须显式分配**, + 留空就意味着接受"谁先抢到谁去下单"。 + +Client 侧对这两种没有任何区别:调 `claim`,拿到任务就做,做完提交。 +**不需要知道这个任务原来有没有主。** 因此 Client 完全不需要看到全局任务池,本地也不缓存未领取的任务 (见 [03 数据模型](03-data-model.md) §3.1)。操作人员想知道队列里还有多少活, @@ -353,7 +372,8 @@ Admin 还没做完,下面这些要和 Admin 一起定。**不许因此停工** **已定案(不再是待确认项):** -- Admin 只把任务分配给指定 Client,Client 不读全局任务池。见 §5.1。 +- 任务分「指定分配」和「无主」两种,`claim` 两种都返回,指定的优先。 + 采集任务不指定,采购任务可指定可留空。Client 不读全局任务池。见 §5.1。 - **没有租约、没有心跳、没有状态回查。** 任务流程只有 §5 / §6 / §7 三个调用;设置页另有 §4.1 的幂等登记调用。 - Admin 必须无条件接受已派发过的 Client 提交的结果。见 §6.1。