Files
cmautobuy/admin/service/submit_test.go
T

512 lines
18 KiB
Go

package service
import (
"database/sql"
"encoding/json"
"errors"
"strings"
"sync"
"testing"
"time"
"cmautobuy/admin/model"
"cmautobuy/admin/repository"
)
// claimTask 走一遍完整的领取流程,让 task_claims 里留下记录。
func claimTask(t *testing.T, db *sql.DB, taskID, clientID string) {
t.Helper()
if err := RegisterClient(db, model.Client{ClientID: clientID}, true); err != nil {
t.Fatalf("注册客户端失败: %v", err)
}
task, err := ClaimNextTask(db, clientID, []string{"collect", "purchase"})
if err != nil {
t.Fatalf("领取失败: %v", err)
}
if task == nil || task.TaskID != taskID {
t.Fatalf("期望领到 %s,实际 %v", taskID, task)
}
}
func taskStatus(t *testing.T, db *sql.DB, taskID string) string {
t.Helper()
var s string
if err := db.QueryRow(`SELECT status FROM tasks WHERE task_id = ?`, taskID).Scan(&s); err != nil {
t.Fatalf("查询任务状态失败: %v", err)
}
return s
}
const resultBody = `{"task_version":1,"attempt_id":"a-1","result_type":"purchase",
"completed_at":"2026-08-06T08:03:00Z","pdd_data":{"schema_version":1,"goods":{"goods_id":"1"}}}`
// ── 正常路径 ───────────────────────────────────────────
func TestSubmitResult_成功落库并标记succeeded(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
resp, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(resultBody))
if err != nil {
t.Fatalf("提交失败: %v", err)
}
var got map[string]any
if err := json.Unmarshal([]byte(resp), &got); err != nil {
t.Fatalf("响应不是合法 JSON: %v", err)
}
if got["accepted"] != true {
t.Errorf("accepted 应为 true,实际 %v", got["accepted"])
}
if got["result_id"] == "" || got["result_id"] == nil {
t.Error("result_id 不应为空")
}
if s := taskStatus(t, db, "TASK-A"); s != "succeeded" {
t.Errorf("任务状态应为 succeeded,实际 %s", s)
}
}
// ── 幂等 ───────────────────────────────────────────────
func TestSubmitResult_同键同内容返回同一结果(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
first, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(resultBody))
if err != nil {
t.Fatalf("首次提交失败: %v", err)
}
second, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(resultBody))
if err != nil {
t.Fatalf("重复提交失败: %v", err)
}
if first != second {
t.Errorf("重复提交应返回完全相同的响应:\n第一次 %s\n第二次 %s", first, second)
}
// 而且不能重复落库
var n int
db.QueryRow(`SELECT COUNT(*) FROM idempotency_keys`).Scan(&n)
if n != 1 {
t.Errorf("幂等表应只有 1 条记录,实际 %d", n)
}
}
func TestSubmitResult_同键不同内容返回冲突(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
if _, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(resultBody)); err != nil {
t.Fatalf("首次提交失败: %v", err)
}
other := `{"task_version":1,"attempt_id":"a-1","result_type":"purchase","pdd_data":{"不一样":true}}`
_, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(other))
if !errors.Is(err, repository.ErrIdempotencyConflict) {
t.Errorf("期望幂等冲突,实际 %v", err)
}
}
// 并发重复提交:只能落库一次。
func TestSubmitResult_并发同键只处理一次(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
const workers = 6
var wg sync.WaitGroup
var mu sync.Mutex
responses := map[string]int{}
var lastErr error
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
resp, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(resultBody))
mu.Lock()
defer mu.Unlock()
if err != nil {
lastErr = err
return
}
responses[resp]++
}()
}
wg.Wait()
if lastErr != nil {
t.Fatalf("并发提交出错: %v", lastErr)
}
if len(responses) != 1 {
t.Errorf("并发提交应返回同一个响应,实际出现 %d 种", len(responses))
}
var n int
db.QueryRow(`SELECT COUNT(*) FROM idempotency_keys`).Scan(&n)
if n != 1 {
t.Errorf("幂等表应只有 1 条记录,实际 %d", n)
}
}
// ── 无条件接受(本工单最要紧的三条)───────────────────
func TestSubmitResult_任务已取消仍然接受(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
// Admin 侧把任务取消了
if _, err := db.Exec(`UPDATE tasks SET status='cancelled' WHERE task_id='TASK-A'`); err != nil {
t.Fatal(err)
}
// 客户端并不知道,照样提交——必须被接受,
// 因为它可能真的已经在拼多多下过单了,这些数据必须留痕
if _, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(resultBody)); err != nil {
t.Fatalf("任务已取消时提交被拒绝了,这违反契约 §4.1: %v", err)
}
if s := taskStatus(t, db, "TASK-A"); s != "succeeded" {
t.Errorf("结果应被记录,状态应为 succeeded,实际 %s", s)
}
}
func TestSubmitResult_任务已重派仍然接受原客户端的结果(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
// 操作员把任务改派给了另一台
if _, err := db.Exec(
`UPDATE tasks SET assigned_client='client-002', status='assigned' WHERE task_id='TASK-A'`,
); err != nil {
t.Fatal(err)
}
// 原来那台还是把结果交上来了 —— 必须接受
if _, err := SubmitResult(db, "TASK-A", "client-001", "key-1", []byte(resultBody)); err != nil {
t.Fatalf("任务已重派时拒绝了原客户端,这违反契约 §4.1: %v", err)
}
}
func TestSubmitResult_同一任务接受多个客户端的多份结果(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
// 重派给第二台,让它也领一次
if _, err := db.Exec(
`UPDATE tasks SET assigned_client='client-002', status='assigned' WHERE task_id='TASK-A'`,
); err != nil {
t.Fatal(err)
}
claimTask(t, db, "TASK-A", "client-002")
// 两台都提交,用各自的幂等键
if _, err := SubmitResult(db, "TASK-A", "client-001", "key-c1", []byte(resultBody)); err != nil {
t.Fatalf("client-001 提交被拒: %v", err)
}
if _, err := SubmitResult(db, "TASK-A", "client-002", "key-c2", []byte(resultBody)); err != nil {
t.Fatalf("client-002 提交被拒: %v", err)
}
var n int
db.QueryRow(`SELECT COUNT(*) FROM idempotency_keys`).Scan(&n)
if n != 2 {
t.Errorf("两份结果应各留一条幂等记录,实际 %d", n)
}
}
// ── 唯一该拒绝的情况 ───────────────────────────────────
func TestSubmitResult_从没领过的客户端被拒绝(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
_, err := SubmitResult(db, "TASK-A", "client-999", "key-x", []byte(resultBody))
if !errors.Is(err, ErrNeverClaimed) {
t.Errorf("从没领过的客户端应被拒绝,实际 %v", err)
}
}
func TestSubmitResult_任务不存在(t *testing.T) {
db := newTestDB(t)
_, err := SubmitResult(db, "TASK-NOPE", "client-001", "key-x", []byte(resultBody))
if !errors.Is(err, ErrTaskNotFound) {
t.Errorf("期望任务不存在错误,实际 %v", err)
}
}
// ── 采集任务的结果要落到 PDD 商品 ──────────────────────
// insertPddProduct 造一个待采集的 PDD 商品。
func insertPddProduct(t *testing.T, db *sql.DB, pddGoodsID string) {
t.Helper()
now := model.NowISO()
if _, err := db.Exec(`
INSERT INTO pdd_products (goods_id, url, collect_status, created_at, updated_at)
VALUES (?, ?, 'collecting', ?, ?)`,
pddGoodsID, "https://mobile.yangkeduo.com/goods.html?goods_id="+pddGoodsID,
now, now); err != nil {
t.Fatalf("插入 PDD 商品失败: %v", err)
}
}
// insertCollectTask 造一个采集任务,关联指定的蝦皮商品和 PDD 商品。
func insertCollectTask(t *testing.T, db *sql.DB, taskID, client, shopeeGoodsID, pddGoodsID string) {
t.Helper()
now := model.NowISO()
if _, err := db.Exec(`
INSERT INTO tasks (task_id, task_type, status, assigned_client,
goods_id, pdd_goods_id, pdd_goods_url, created_at, updated_at)
VALUES (?, 'collect', 'assigned', ?, ?, ?, 'https://x/1', ?, ?)`,
taskID, client, shopeeGoodsID, pddGoodsID, now, now); err != nil {
t.Fatalf("插入采集任务失败: %v", err)
}
}
// pddProductState 读回 PDD 商品的采集状态、标题和结果。
func pddProductState(t *testing.T, db *sql.DB, goodsID string) (status, title, skus, msg string) {
t.Helper()
var ti, sk, ms sql.NullString
if err := db.QueryRow(`
SELECT collect_status, title, skus_json, collect_msg
FROM pdd_products WHERE goods_id = ?`, goodsID,
).Scan(&status, &ti, &sk, &ms); err != nil {
t.Fatalf("查询 PDD 商品失败: %v", err)
}
return status, ti.String, sk.String, ms.String
}
const collectBody = `{"task_version":1,"attempt_id":"a-1","result_type":"collect",
"pdd_data":{"schema_version":1,"goods_id":"PDD-1","title":"测试商品",
"shop_name":"测试旗舰店","price_granularity":"color",
"dimensions":[{"key":"color","name":"颜色分类"}],
"skus":[{"options":{"color":"黑色","size":"M"},"price_cent":1256,
"price_observed_at":{"color":"黑色","size":"M"},
"available":true,"raw_price":"¥12.56"}]}}`
func TestSubmitResult_采集结果写入PDD商品(t *testing.T) {
db := newTestDB(t)
insertPddProduct(t, db, "PDD-1")
insertCollectTask(t, db, "TASK-C", "client-001", "SHOPEE-1", "PDD-1")
claimTask(t, db, "TASK-C", "client-001")
if _, err := SubmitResult(db, "TASK-C", "client-001", "key-c", []byte(collectBody)); err != nil {
t.Fatalf("提交采集结果失败: %v", err)
}
status, title, skus, _ := pddProductState(t, db, "PDD-1")
if status != "collected" {
t.Errorf("采集状态应为 collected,实际 %s", status)
}
if title != "测试商品" {
t.Errorf("标题应从采集结果里取出来,实际 %q", title)
}
if skus == "" {
t.Error("skus_json 应被写入,实际为空")
}
var shop sql.NullString
if err := db.QueryRow(
`SELECT shop_name FROM pdd_products WHERE goods_id = ?`, "PDD-1",
).Scan(&shop); err != nil {
t.Fatalf("读取店铺名失败: %v", err)
}
if !shop.Valid || shop.String != "测试旗舰店" {
t.Errorf("店铺名应从采集结果落库,实际 %+v", shop)
}
if !strings.Contains(skus, `"price_granularity":"color"`) ||
!strings.Contains(skus, `"price_observed_at"`) {
t.Errorf("价格粒度与实测组合应保留在 skus_json,实际 %s", skus)
}
}
func TestSubmitResult_没有店铺名时不覆盖已有值(t *testing.T) {
db := newTestDB(t)
insertPddProduct(t, db, "PDD-1")
if _, err := db.Exec(
`UPDATE pdd_products SET shop_name = '旧店铺' WHERE goods_id = 'PDD-1'`,
); err != nil {
t.Fatalf("准备已有店铺名失败: %v", err)
}
insertCollectTask(t, db, "TASK-C", "client-001", "SHOPEE-1", "PDD-1")
claimTask(t, db, "TASK-C", "client-001")
body := `{"attempt_id":"a-1","result_type":"collect",
"pdd_data":{"goods_id":"PDD-1","title":"测试商品",
"skus":[{"options":{"color":"黑色"},"price_cent":100,"available":true}]}}`
if _, err := SubmitResult(db, "TASK-C", "client-001", "key-no-shop", []byte(body)); err != nil {
t.Fatalf("老版本报文不应失败: %v", err)
}
var shop string
if err := db.QueryRow(
`SELECT shop_name FROM pdd_products WHERE goods_id = 'PDD-1'`,
).Scan(&shop); err != nil {
t.Fatalf("读取店铺名失败: %v", err)
}
if shop != "旧店铺" {
t.Errorf("缺少 shop_name 时不应覆盖已有值,实际 %q", shop)
}
}
func TestSubmitResult_拒绝未知价格采样粒度(t *testing.T) {
db := newTestDB(t)
insertPddProduct(t, db, "PDD-1")
insertCollectTask(t, db, "TASK-C", "client-001", "SHOPEE-1", "PDD-1")
claimTask(t, db, "TASK-C", "client-001")
body := `{"attempt_id":"a-1","result_type":"collect",
"pdd_data":{"goods_id":"PDD-1","price_granularity":"unknown",
"skus":[{"options":{"color":"黑色"},"price_cent":100,"available":true}]}}`
if _, err := SubmitResult(db, "TASK-C", "client-001", "key-bad-granularity", []byte(body)); err == nil {
t.Fatal("未知 price_granularity 应被拒绝")
}
}
// 采到 0 个规格必须算失败——数据对业务毫无用处,
// 显示"已采集"会让操作员以为好了,等建任务时才发现不对。
func TestSubmitResult_采到零个规格算失败(t *testing.T) {
db := newTestDB(t)
insertPddProduct(t, db, "PDD-1")
insertCollectTask(t, db, "TASK-C", "client-001", "SHOPEE-1", "PDD-1")
claimTask(t, db, "TASK-C", "client-001")
body := `{"attempt_id":"a-1","result_type":"collect",
"pdd_data":{"schema_version":1,"goods_id":"PDD-1","skus":[]}}`
if _, err := SubmitResult(db, "TASK-C", "client-001", "key-c", []byte(body)); err != nil {
t.Fatalf("提交不该报错,应该记成采集失败: %v", err)
}
status, _, _, msg := pddProductState(t, db, "PDD-1")
if status != "failed" {
t.Errorf("采到 0 个规格应记为 failed,实际 %s", status)
}
if msg == "" {
t.Error("应写明失败原因,否则操作员不知道为什么")
}
}
// 客户端采回来的商品和请求的对不上时必须拒绝——
// 不拦就会把 B 的规格价格存到 A 名下,之后按它下单就是买错东西。
func TestSubmitResult_采错商品被拒绝(t *testing.T) {
db := newTestDB(t)
insertPddProduct(t, db, "PDD-1")
insertCollectTask(t, db, "TASK-C", "client-001", "SHOPEE-1", "PDD-1")
claimTask(t, db, "TASK-C", "client-001")
// 请求采 PDD-1,客户端却返回了 PDD-999
body := `{"attempt_id":"a-1","result_type":"collect",
"pdd_data":{"schema_version":1,"goods_id":"PDD-999","title":"别的商品",
"skus":[{"options":{"color":"黑色"},"price_cent":100,"available":true}]}}`
_, err := SubmitResult(db, "TASK-C", "client-001", "key-c", []byte(body))
if !errors.Is(err, ErrCollectMismatch) {
t.Fatalf("期望 ErrCollectMismatch,实际 %v", err)
}
// 事务应整体回滚,PDD 商品不能被写脏
status, _, skus, _ := pddProductState(t, db, "PDD-1")
if status == "collected" || skus != "" {
t.Errorf("采错商品时不该写入任何结果,实际 status=%s skus=%q", status, skus)
}
}
// ── 提交失败 ───────────────────────────────────────────
func TestSubmitFailure_状态映射(t *testing.T) {
cases := []struct {
reported string
want string
}{
{"retry_wait", "assigned"}, // 放回去等它再来领
{"manual_review", "manual_review"},
{"failed", "failed"},
{"cancelled", "cancelled"},
}
for _, tc := range cases {
t.Run(tc.reported, func(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
body := `{"task_version":1,"attempt_id":"a-1","status":"` + tc.reported +
`","error":{"code":"PDD_PAGE_TIMEOUT","message":"页面超时"}}`
response, err := SubmitFailure(db, "TASK-A", "client-001", "k", []byte(body))
if err != nil {
t.Fatalf("提交失败结果出错: %v", err)
}
var receipt map[string]any
if err := json.Unmarshal([]byte(response), &receipt); err != nil {
t.Fatalf("失败响应不是 JSON: %v", err)
}
if receipt["result_id"] == "" || receipt["result_id"] == nil {
t.Error("失败响应 result_id 不应为空")
}
if s := taskStatus(t, db, "TASK-A"); s != tc.want {
t.Errorf("%s 应映射为 %s,实际 %s", tc.reported, tc.want, s)
}
})
}
}
func TestSubmitFailure_非法状态被拒绝(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
body := `{"attempt_id":"a-1","status":"succeeded"}`
if _, err := SubmitFailure(db, "TASK-A", "client-001", "k", []byte(body)); err == nil {
t.Error("失败接口不该接受 succeeded 这种状态")
}
}
func TestSubmitFailure_采集失败写回PDD商品(t *testing.T) {
db := newTestDB(t)
insertPddProduct(t, db, "PDD-1")
insertCollectTask(t, db, "TASK-C", "client-001", "SHOPEE-1", "PDD-1")
claimTask(t, db, "TASK-C", "client-001")
body := `{"attempt_id":"a-1","status":"failed",
"error":{"code":"PDD_PAGE_TIMEOUT","message":"商品页加载超时"}}`
if _, err := SubmitFailure(db, "TASK-C", "client-001", "k", []byte(body)); err != nil {
t.Fatalf("提交失败: %v", err)
}
status, _, _, msg := pddProductState(t, db, "PDD-1")
if status != "failed" {
t.Errorf("采集状态应为 failed,实际 %s", status)
}
if msg == "" {
t.Error("失败原因应写入 collect_msg,否则操作员看不到为什么采不到")
}
}
// ── 活动时间 ───────────────────────────────────────────
func TestSubmit_刷新客户端活动时间(t *testing.T) {
db := newTestDB(t)
insertTask(t, db, "TASK-A", "client-001")
claimTask(t, db, "TASK-A", "client-001")
// 把客户端改成很久没活动
db.Exec(`UPDATE clients SET last_seen_at='2020-01-01T00:00:00Z' WHERE client_id='client-001'`)
views, _ := ListClientViews(db, "", 10*time.Minute)
if views[0].Status != "离线" {
t.Fatalf("前置条件不对,应为离线")
}
if _, err := SubmitResult(db, "TASK-A", "client-001", "k", []byte(resultBody)); err != nil {
t.Fatalf("提交失败: %v", err)
}
views, _ = ListClientViews(db, "", 10*time.Minute)
if views[0].Status != "在线" {
t.Error("提交结果后应刷新活动时间——只在 claim 里刷新的话," +
"客户端执行长任务期间会被误判成离线")
}
}