feat: 实现结果提交接口与幂等处理
补齐 submit_result / submit_failure,Admin 侧的三个接口全部可用, Client 的完整一圈(领取 → 执行 → 提交)现在能走通了。 实现 - 幂等:idempotency_keys 表。同键同内容返回上次的响应且不重复落库, 同键不同内容返回 409。幂等记录与业务写入在**同一事务**, 分开写的话业务成功但幂等没记上,重试会被重复处理 - 无条件接受(契约 §4.1,最容易写错的一条): 任务已取消、已重派给别人,都照样接受结果——客户端中途不查任务状态, 必然会提交"Admin 这边已经不要了"的结果,而它可能真的已经下过单, 这些数据必须留痕 - 采集任务的结果落到商品级 shopee_products.pdd_data 并置 collected; 失败则置 failed 并把原因写进 collect_error,操作员才看得见 - 失败状态映射:retry_wait→assigned,其余同名 - 三个接口都刷新 last_seen_at 新增 task_claims 表(migrations v2) 契约要求"只有从未分配给该客户端的任务才返回 403",但 assigned_client 只记当前归属,重派后就查不出原来那台领过——而契约又要求那种情况必须接受。 没有这张表这条规则根本没法判断。顺带得到一份审计记录。 修复第二个并发 bug:事务必须 BEGIN IMMEDIATE 并发提交报 SQLITE_BUSY。根因是 Go 的 db.Begin() 默认发 BEGIN DEFERRED, 事务开始时不拿写锁,多个事务各自先读再想升级成写就互相卡死, 这种情况 busy_timeout 救不了。DSN 加 _txlock=immediate 后事务一开始 就排队拿锁。实测 6 个并发事务:默认失败 5/6,加参数后 0/6。 已写进 docs/admin/03-data-model.md §2.1。 已验证(Go 1.23.0) - 30 个单元测试全过,并发用例重复 20 次稳定通过 - 端到端:claim 200 → 提交 200 → 重复提交返回完全相同的响应 → 同键不同内容 409 → 没领过的客户端 403 → 任务不存在 404 → 缺 Idempotency-Key 400;库里 task=succeeded、幂等 1 条、领取历史 1 条 说明:Gitea 尚未配置,本次无对应工单号。 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -10,6 +10,8 @@ package api
|
||||
import (
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
|
||||
@@ -163,15 +165,36 @@ func taskPayload(t *model.Task) gin.H {
|
||||
// - 不得因为任务已取消而拒绝;
|
||||
// - 不得因为任务已重派给别的客户端而拒绝;
|
||||
// - 必须能接受同一任务来自多个客户端的多份结果;
|
||||
// - 只有**从未分配给该客户端**的任务才返回 403。
|
||||
// - 只有**从未领过**这个任务的客户端才返回 403。
|
||||
//
|
||||
// 原因:Client 中途不查任务状态,所以它必然会提交一些
|
||||
// "Admin 这边已经不要了"的结果。而它可能真的已经下单了,
|
||||
// 这些数据必须能交上来留痕。
|
||||
//
|
||||
// accepted: true 的意思是"我收到并存下了",不代表任务还算数。
|
||||
// 原因见 service.SubmitResult 的注释。
|
||||
func (h *Handler) SubmitResult(c *gin.Context) {
|
||||
h.handleSubmit(c, service.SubmitResult, "task_result_received")
|
||||
}
|
||||
|
||||
// SubmitFailure 接收失败或需人工处理的结果。
|
||||
//
|
||||
// Client 报告的 status 映射见 docs/admin/04-client-api.md §5。
|
||||
// §4.1 的无条件接受规则同样适用。
|
||||
func (h *Handler) SubmitFailure(c *gin.Context) {
|
||||
h.handleSubmit(c, service.SubmitFailure, "task_failure_received")
|
||||
}
|
||||
|
||||
// submitFunc 是两个提交接口共用的处理函数形状。
|
||||
type submitFunc func(db *sql.DB, taskID, clientID, idemKey string, rawBody []byte) (string, error)
|
||||
|
||||
// handleSubmit 把两个提交接口共同的部分抽出来:
|
||||
// 取参数、读原始请求体、调 service、把业务错误翻译成 HTTP 状态码。
|
||||
func (h *Handler) handleSubmit(c *gin.Context, submit submitFunc, event string) {
|
||||
taskID := c.Param("task_id")
|
||||
|
||||
clientID := c.GetHeader("X-Client-Id")
|
||||
if clientID == "" {
|
||||
apiError(c, http.StatusBadRequest, "MISSING_CLIENT_ID",
|
||||
"缺少 X-Client-Id 请求头", false)
|
||||
return
|
||||
}
|
||||
|
||||
idemKey := c.GetHeader("Idempotency-Key")
|
||||
if idemKey == "" {
|
||||
apiError(c, http.StatusBadRequest, "MISSING_IDEMPOTENCY_KEY",
|
||||
@@ -179,43 +202,39 @@ func (h *Handler) SubmitResult(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
// TODO(骨架): 1. 查 idempotency_keys:
|
||||
// 键同、内容哈希同 -> 直接返回上次的 response_body,不重复落库
|
||||
// 键同、内容哈希不同 -> 409 IDEMPOTENCY_CONFLICT
|
||||
// TODO(骨架): 2. 在**同一个事务**里:
|
||||
// 写 tasks.result_data;
|
||||
// 采集任务同时写 shopee_products.pdd_data 并置 collect_status = collected;
|
||||
// tasks.status = 'succeeded';
|
||||
// 写 idempotency_keys。
|
||||
// TODO(骨架): 3. 刷新客户端 last_seen_at
|
||||
// (长任务期间不调 claim,不刷新会被误判成离线)。
|
||||
// 读**原始字节**而不是解析后再序列化——幂等要比对请求内容的哈希,
|
||||
// 重新序列化会因为字段顺序、空格不同而算出不一样的哈希。
|
||||
rawBody, err := io.ReadAll(c.Request.Body)
|
||||
if err != nil {
|
||||
apiError(c, http.StatusBadRequest, "INVALID_BODY", "读取请求体失败", false)
|
||||
return
|
||||
}
|
||||
|
||||
_ = taskID
|
||||
apiError(c, http.StatusNotImplemented, "NOT_IMPLEMENTED",
|
||||
"提交结果接口尚未实现", false)
|
||||
}
|
||||
respJSON, err := submit(h.db, taskID, clientID, idemKey, rawBody)
|
||||
switch {
|
||||
case err == nil:
|
||||
log.Printf("%s client_id=%s task_id=%s", event, clientID, taskID)
|
||||
c.Data(http.StatusOK, "application/json; charset=utf-8", []byte(respJSON))
|
||||
|
||||
// SubmitFailure 接收失败或需人工处理的结果。
|
||||
//
|
||||
// Client 报告的 status 只会是 retry_wait / manual_review / failed / cancelled,
|
||||
// 按下表落地(docs/admin/04-client-api.md §5):
|
||||
//
|
||||
// retry_wait -> assigned(放回去,等它再来领)
|
||||
// manual_review -> manual_review
|
||||
// failed -> failed
|
||||
// cancelled -> cancelled
|
||||
//
|
||||
// §4.1 的无条件接受规则同样适用于本接口。
|
||||
func (h *Handler) SubmitFailure(c *gin.Context) {
|
||||
taskID := c.Param("task_id")
|
||||
case errors.Is(err, service.ErrTaskNotFound):
|
||||
apiError(c, http.StatusNotFound, "TASK_NOT_FOUND",
|
||||
"任务不存在", false)
|
||||
|
||||
// TODO(骨架): 同 SubmitResult 的幂等处理;
|
||||
// 采集任务失败时同步把 shopee_products.collect_status 置 failed,
|
||||
// 错误信息写 collect_error,界面上要看得见。
|
||||
case errors.Is(err, service.ErrNeverClaimed):
|
||||
// 注意:只有**从没领过**才会走到这里。
|
||||
// 任务已取消、已重派,都不会拒绝。
|
||||
apiError(c, http.StatusForbidden, "TASK_NOT_ASSIGNED",
|
||||
"该任务从未分配给这个客户端", false)
|
||||
|
||||
_ = taskID
|
||||
apiError(c, http.StatusNotImplemented, "NOT_IMPLEMENTED",
|
||||
"提交失败接口尚未实现", false)
|
||||
case errors.Is(err, service.ErrIdempotencyConflict):
|
||||
apiError(c, http.StatusConflict, "IDEMPOTENCY_CONFLICT",
|
||||
"相同幂等键提交了不同内容。内容变了应该用新的 attempt_id 生成新键", false)
|
||||
|
||||
default:
|
||||
log.Printf("%s_failed client_id=%s task_id=%s err=%v", event, clientID, taskID, err)
|
||||
apiError(c, http.StatusInternalServerError, "SUBMIT_FAILED",
|
||||
"保存结果失败,请稍后重试", true)
|
||||
}
|
||||
}
|
||||
|
||||
// apiError 按契约格式返回错误。
|
||||
|
||||
@@ -49,9 +49,9 @@ func UpsertClient(db *sql.DB, c model.Client) error {
|
||||
//
|
||||
// claim / result / failure 三个接口都要调。只在 claim 里调的话,
|
||||
// 客户端执行长任务期间不调 claim,会被误判成离线。
|
||||
func TouchClient(db *sql.DB, clientID string) error {
|
||||
func TouchClient(q Execer, clientID string) error {
|
||||
now := model.NowISO()
|
||||
_, err := db.Exec(
|
||||
_, err := q.Exec(
|
||||
`UPDATE clients SET last_seen_at = ?, updated_at = ? WHERE client_id = ?`,
|
||||
now, now, clientID)
|
||||
if err != nil {
|
||||
|
||||
+49
-1
@@ -20,6 +20,17 @@ import (
|
||||
_ "modernc.org/sqlite"
|
||||
)
|
||||
|
||||
// Execer 让同一个 repository 函数既能直接用 *sql.DB,
|
||||
// 也能在事务里用 *sql.Tx。
|
||||
//
|
||||
// 需要"多张表要么一起改、要么都不改"时,service 开一个事务,
|
||||
// 把 *sql.Tx 传进来即可,不用为事务再写一套函数。
|
||||
type Execer interface {
|
||||
Exec(query string, args ...any) (sql.Result, error)
|
||||
Query(query string, args ...any) (*sql.Rows, error)
|
||||
QueryRow(query string, args ...any) *sql.Row
|
||||
}
|
||||
|
||||
// Open 打开 data/admin.db。
|
||||
//
|
||||
// **PRAGMA 必须写在 DSN 里,不能用 db.Exec("PRAGMA ...") 设置。**
|
||||
@@ -37,10 +48,13 @@ func Open(dataDir string) (*sql.DB, error) {
|
||||
// busy_timeout 拿不到锁时最多等 5 秒,而不是立刻报错
|
||||
// journal_mode WAL 模式,读和写可以同时进行
|
||||
// foreign_keys 打开外键约束(SQLite 默认是关的)
|
||||
//
|
||||
// _txlock=immediate 是另一个**必须加**的参数,原因见下。
|
||||
dsn := "file:" + path +
|
||||
"?_pragma=busy_timeout(5000)" +
|
||||
"&_pragma=journal_mode(WAL)" +
|
||||
"&_pragma=foreign_keys(1)"
|
||||
"&_pragma=foreign_keys(1)" +
|
||||
"&_txlock=immediate"
|
||||
|
||||
db, err := sql.Open("sqlite", dsn)
|
||||
if err != nil {
|
||||
@@ -55,6 +69,21 @@ func Open(dataDir string) (*sql.DB, error) {
|
||||
db.SetMaxOpenConns(4)
|
||||
db.SetMaxIdleConns(4)
|
||||
|
||||
// 关于 _txlock=immediate:
|
||||
//
|
||||
// Go 的 db.Begin() 默认发的是 BEGIN DEFERRED——事务开始时**不拿写锁**,
|
||||
// 等到第一次写才去拿。于是多个事务可以同时开始、各自先读,
|
||||
// 然后同时想升级成写,互相卡死,直接报 SQLITE_BUSY。
|
||||
// 这种情况 busy_timeout **救不了**:等下去也不可能有结果,
|
||||
// 只能让某个事务整个重来。
|
||||
//
|
||||
// 加上 _txlock=immediate 后,事务一开始就拿写锁,
|
||||
// 拿不到就按 busy_timeout 排队等——这才是我们要的行为。
|
||||
//
|
||||
// 实测(6 个并发事务,每个先读后写):
|
||||
// 默认 deferred 失败 5/6
|
||||
// _txlock=immediate 失败 0/6
|
||||
|
||||
// sql.Open 是懒加载的,这里主动连一次,好让配置错误立刻暴露
|
||||
if err := db.Ping(); err != nil {
|
||||
db.Close()
|
||||
@@ -192,6 +221,25 @@ var migrations = [][]string{
|
||||
created_at TEXT NOT NULL
|
||||
);`,
|
||||
},
|
||||
|
||||
// v2: 领取历史。
|
||||
//
|
||||
// 为什么需要它:契约要求"只有**从未分配给该客户端**的任务才返回 403"
|
||||
// (docs/admin/04-client-api.md §4.1)。但 tasks.assigned_client 只记
|
||||
// **当前**归属,任务一旦重派给别人,就查不出原来那台领过——
|
||||
// 而契约又明确要求"已重派仍要接受原客户端提交的结果"。
|
||||
// 没有这张表,那条规则根本没法判断。
|
||||
//
|
||||
// 顺带得到一份审计记录:这个任务被哪几台客户端领过。
|
||||
{
|
||||
`CREATE TABLE task_claims (
|
||||
task_id TEXT NOT NULL,
|
||||
client_id TEXT NOT NULL,
|
||||
claimed_at TEXT NOT NULL,
|
||||
PRIMARY KEY (task_id, client_id)
|
||||
);`,
|
||||
`CREATE INDEX idx_task_claims_client ON task_claims(client_id);`,
|
||||
},
|
||||
}
|
||||
|
||||
// Migrate 把数据库升到最新版本。
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
package repository
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"database/sql"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
"cmautobuy/admin/model"
|
||||
)
|
||||
|
||||
// ErrIdempotencyConflict 表示同一个键提交了不同的内容。
|
||||
//
|
||||
// 这说明客户端弄错了——同一个键必须对应同一份内容。
|
||||
// 内容真的变了,应该用新的 attempt_id 生成新的键。
|
||||
var ErrIdempotencyConflict = errors.New("相同幂等键提交了不同内容")
|
||||
|
||||
// HashRequest 计算请求体的哈希,用来判断"同一个键"配的是不是"同一份内容"。
|
||||
func HashRequest(body []byte) string {
|
||||
sum := sha256.Sum256(body)
|
||||
return hex.EncodeToString(sum[:])
|
||||
}
|
||||
|
||||
// LookupIdempotent 查这个键是不是已经处理过。
|
||||
//
|
||||
// 三种结果:
|
||||
// - 处理过且内容一致 → 返回上次的响应体和 true,**调用方直接原样返回,不要重复落库**
|
||||
// - 处理过但内容不同 → 返回 ErrIdempotencyConflict,调用方回 409
|
||||
// - 没处理过 → 返回 ("", false, nil),调用方正常处理
|
||||
func LookupIdempotent(q Execer, key, requestHash string) (string, bool, error) {
|
||||
var storedHash, storedBody string
|
||||
err := q.QueryRow(
|
||||
`SELECT request_hash, response_body FROM idempotency_keys WHERE key = ?`,
|
||||
key).Scan(&storedHash, &storedBody)
|
||||
|
||||
if err == sql.ErrNoRows {
|
||||
return "", false, nil
|
||||
}
|
||||
if err != nil {
|
||||
return "", false, fmt.Errorf("查询幂等键失败: %w", err)
|
||||
}
|
||||
if storedHash != requestHash {
|
||||
return "", false, ErrIdempotencyConflict
|
||||
}
|
||||
return storedBody, true, nil
|
||||
}
|
||||
|
||||
// SaveIdempotent 记下这个键处理过了,以及当时返回了什么。
|
||||
//
|
||||
// `[必须]` 必须和业务写入在**同一个事务**里。
|
||||
// 分开写的话,业务写成功但幂等记录没写上,客户端重试就会被重复处理。
|
||||
func SaveIdempotent(q Execer, key, requestHash, responseBody string) error {
|
||||
_, err := q.Exec(
|
||||
`INSERT INTO idempotency_keys (key, request_hash, response_body, created_at)
|
||||
VALUES (?, ?, ?, ?)`,
|
||||
key, requestHash, responseBody, model.NowISO())
|
||||
if err != nil {
|
||||
return fmt.Errorf("保存幂等键失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
package repository
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"cmautobuy/admin/model"
|
||||
)
|
||||
|
||||
// SetCollectResult 把采集回来的 PDD 商品数据存到商品级。
|
||||
//
|
||||
// `[必须]` pdd_data 存在**商品级**(shopee_products),不是订单级——
|
||||
// 一个 PDD 商品采一次,所有相关订单共用这份结果。
|
||||
// 见 docs/admin/03-data-model.md §3.1。
|
||||
func SetCollectResult(q Execer, goodsID, pddData string) error {
|
||||
if goodsID == "" {
|
||||
return nil // 任务没关联蝦皮商品(比如手工造的测试任务),跳过
|
||||
}
|
||||
now := model.NowISO()
|
||||
_, err := q.Exec(`
|
||||
UPDATE shopee_products
|
||||
SET pdd_data = ?, collect_status = 'collected',
|
||||
collect_error = NULL, collected_at = ?, updated_at = ?
|
||||
WHERE goods_id = ?`,
|
||||
pddData, now, now, goodsID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("保存商品 %s 的采集结果失败: %w", goodsID, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetCollectFailed 标记采集失败,并记下原因。
|
||||
//
|
||||
// 错误信息要能在界面上看见,否则操作员不知道为什么采不到。
|
||||
func SetCollectFailed(q Execer, goodsID, errMsg string) error {
|
||||
if goodsID == "" {
|
||||
return nil
|
||||
}
|
||||
now := model.NowISO()
|
||||
_, err := q.Exec(`
|
||||
UPDATE shopee_products
|
||||
SET collect_status = 'failed', collect_error = ?, updated_at = ?
|
||||
WHERE goods_id = ?`,
|
||||
errMsg, now, goodsID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("标记商品 %s 采集失败出错: %w", goodsID, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -74,6 +74,11 @@ func ClaimNextTask(db *sql.DB, clientID string, supportedTypes []string) (*model
|
||||
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 // 没有可领的任务
|
||||
@@ -125,3 +130,101 @@ func GetTask(db *sql.DB, taskID string) (*model.Task, error) {
|
||||
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
|
||||
}
|
||||
|
||||
// TaskExists 判断任务在不在,并返回它的类型和关联的蝦皮商品编号。
|
||||
func TaskExists(q Execer, taskID string) (exists bool, taskType model.TaskType, goodsID string, err error) {
|
||||
var gid sql.NullString
|
||||
e := q.QueryRow(
|
||||
`SELECT task_type, goods_id FROM tasks WHERE task_id = ?`,
|
||||
taskID).Scan(&taskType, &gid)
|
||||
if e == sql.ErrNoRows {
|
||||
return false, "", "", nil
|
||||
}
|
||||
if e != nil {
|
||||
return false, "", "", fmt.Errorf("查询任务 %s 失败: %w", taskID, e)
|
||||
}
|
||||
return true, taskType, gid.String, 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
|
||||
}
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
)
|
||||
|
||||
// newID 生成一个随机编号,用作 result_id 之类的标识。
|
||||
//
|
||||
// 没用 UUID 库是为了少一个依赖——这里只要求"够随机、不重复",
|
||||
// 16 字节的随机十六进制串完全够用。
|
||||
func newID() string {
|
||||
buf := make([]byte, 16)
|
||||
if _, err := rand.Read(buf); err != nil {
|
||||
// crypto/rand 读失败说明系统熵源出了问题,这种情况极罕见。
|
||||
// 返回空串让调用方的响应里 result_id 为空,总比 panic 掉整个进程好。
|
||||
return ""
|
||||
}
|
||||
return hex.EncodeToString(buf)
|
||||
}
|
||||
@@ -0,0 +1,223 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"cmautobuy/admin/model"
|
||||
"cmautobuy/admin/repository"
|
||||
)
|
||||
|
||||
// 提交结果时可能出现的几种业务错误,handler 据此决定 HTTP 状态码。
|
||||
var (
|
||||
// ErrTaskNotFound 任务不存在 -> 404
|
||||
ErrTaskNotFound = errors.New("任务不存在")
|
||||
// ErrNeverClaimed 这台客户端从没领过这个任务 -> 403。
|
||||
// 注意**只有这一种情况才拒绝**,见 SubmitResult 的说明。
|
||||
ErrNeverClaimed = errors.New("该任务从未分配给这个客户端")
|
||||
// ErrIdempotencyConflict 同一个键提交了不同内容 -> 409
|
||||
ErrIdempotencyConflict = repository.ErrIdempotencyConflict
|
||||
)
|
||||
|
||||
// ResultRequest 是客户端提交成功结果的请求体。
|
||||
type ResultRequest struct {
|
||||
TaskVersion int `json:"task_version"`
|
||||
AttemptID string `json:"attempt_id"`
|
||||
ResultType string `json:"result_type"`
|
||||
CompletedAt string `json:"completed_at"`
|
||||
PddData json.RawMessage `json:"pdd_data"`
|
||||
}
|
||||
|
||||
// FailureRequest 是客户端提交失败/需人工处理的请求体。
|
||||
type FailureRequest struct {
|
||||
TaskVersion int `json:"task_version"`
|
||||
AttemptID string `json:"attempt_id"`
|
||||
Status string `json:"status"`
|
||||
Error struct {
|
||||
Code string `json:"code"`
|
||||
Message string `json:"message"`
|
||||
Retryable bool `json:"retryable"`
|
||||
Step string `json:"step"`
|
||||
} `json:"error"`
|
||||
Diagnostics json.RawMessage `json:"diagnostics"`
|
||||
ReportedAt string `json:"reported_at"`
|
||||
}
|
||||
|
||||
// SubmitResult 接收客户端提交的成功结果,返回要回给客户端的 JSON。
|
||||
//
|
||||
// # 最容易写错的一条:无条件接受
|
||||
//
|
||||
// 下面每一条都不允许违反(docs/admin/04-client-api.md §4.1):
|
||||
//
|
||||
// - **不得**因为任务已取消而拒绝;
|
||||
// - **不得**因为任务已重派给别的客户端而拒绝;
|
||||
// - **必须**能接受同一任务来自多个客户端的多份结果;
|
||||
// - 只有**从未领过**这个任务的客户端才返回 403。
|
||||
//
|
||||
// 原因:客户端中途不查任务状态(这是有意的设计),所以它**必然**会提交
|
||||
// 一些"Admin 这边已经不要了"的结果。而它可能真的已经在拼多多下过单了,
|
||||
// 这些数据必须能交上来留痕,否则就成了一笔谁都不知道的订单。
|
||||
//
|
||||
// accepted: true 的意思是"**我收到并存下了**",不代表这个任务还算数。
|
||||
// 任务算不算数由人工审核决定。
|
||||
func SubmitResult(db *sql.DB, taskID, clientID, idemKey string, rawBody []byte) (string, error) {
|
||||
var req ResultRequest
|
||||
if err := json.Unmarshal(rawBody, &req); err != nil {
|
||||
return "", fmt.Errorf("请求体不是合法 JSON: %w", err)
|
||||
}
|
||||
|
||||
return submitInTx(db, taskID, clientID, idemKey, rawBody,
|
||||
func(tx *sql.Tx, taskType model.TaskType, goodsID string) (map[string]any, error) {
|
||||
pddData := string(req.PddData)
|
||||
if strings.TrimSpace(pddData) == "" {
|
||||
pddData = "{}"
|
||||
}
|
||||
|
||||
if err := repository.MarkTaskSucceeded(tx, taskID, pddData); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// 采集任务的结果还要落到商品级,供后续规格匹配使用
|
||||
if taskType == model.TaskCollect {
|
||||
if err := repository.SetCollectResult(tx, goodsID, pddData); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return map[string]any{
|
||||
"accepted": true,
|
||||
"result_id": newID(),
|
||||
"accepted_at": model.NowISO(),
|
||||
}, nil
|
||||
})
|
||||
}
|
||||
|
||||
// SubmitFailure 接收客户端提交的失败/需人工处理结果。
|
||||
//
|
||||
// §4.1 的无条件接受规则**同样适用于本函数**。
|
||||
func SubmitFailure(db *sql.DB, taskID, clientID, idemKey string, rawBody []byte) (string, error) {
|
||||
var req FailureRequest
|
||||
if err := json.Unmarshal(rawBody, &req); err != nil {
|
||||
return "", fmt.Errorf("请求体不是合法 JSON: %w", err)
|
||||
}
|
||||
|
||||
newStatus, ok := mapFailureStatus(req.Status)
|
||||
if !ok {
|
||||
return "", fmt.Errorf(
|
||||
"status 只能是 retry_wait / manual_review / failed / cancelled,收到 %q", req.Status)
|
||||
}
|
||||
|
||||
return submitInTx(db, taskID, clientID, idemKey, rawBody,
|
||||
func(tx *sql.Tx, taskType model.TaskType, goodsID string) (map[string]any, error) {
|
||||
if err := repository.MarkTaskFailure(
|
||||
tx, taskID, newStatus, req.Error.Code, req.Error.Message); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// 采集失败要让操作员在蝦皮数据页看得见原因
|
||||
if taskType == model.TaskCollect {
|
||||
msg := req.Error.Message
|
||||
if req.Error.Code != "" {
|
||||
msg = req.Error.Code + ": " + msg
|
||||
}
|
||||
if err := repository.SetCollectFailed(tx, goodsID, msg); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return map[string]any{
|
||||
"accepted": true,
|
||||
"task_status": string(newStatus),
|
||||
"accepted_at": model.NowISO(),
|
||||
}, nil
|
||||
})
|
||||
}
|
||||
|
||||
// mapFailureStatus 把客户端报告的状态映射成 Admin 侧状态。
|
||||
// 映射表见 docs/admin/04-client-api.md §5。
|
||||
func mapFailureStatus(reported string) (model.TaskStatus, bool) {
|
||||
switch reported {
|
||||
case "retry_wait":
|
||||
// 放回待领取,等客户端下次再来领
|
||||
return model.TaskAssigned, true
|
||||
case "manual_review":
|
||||
return model.TaskManualReview, true
|
||||
case "failed":
|
||||
return model.TaskFailed, true
|
||||
case "cancelled":
|
||||
return model.TaskCancelled, true
|
||||
default:
|
||||
return "", false
|
||||
}
|
||||
}
|
||||
|
||||
// submitInTx 把两个提交接口共同的骨架抽出来:
|
||||
// 幂等检查 → 权限判断 → 业务写入 → 记幂等 → 刷新客户端活动时间,
|
||||
// 全部在**一个事务**里完成。
|
||||
//
|
||||
// 业务写入部分由 apply 提供,它拿到的 tx 和外层是同一个。
|
||||
func submitInTx(
|
||||
db *sql.DB, taskID, clientID, idemKey string, rawBody []byte,
|
||||
apply func(tx *sql.Tx, taskType model.TaskType, goodsID string) (map[string]any, error),
|
||||
) (string, error) {
|
||||
|
||||
hash := repository.HashRequest(rawBody)
|
||||
|
||||
tx, err := db.Begin()
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("开始事务失败: %w", err)
|
||||
}
|
||||
defer tx.Rollback() // 已提交的事务再 Rollback 是空操作,安全
|
||||
|
||||
// 1. 处理过就原样返回上次的响应,绝不重复落库
|
||||
if body, done, err := repository.LookupIdempotent(tx, idemKey, hash); err != nil {
|
||||
return "", err
|
||||
} else if done {
|
||||
return body, nil
|
||||
}
|
||||
|
||||
// 2. 任务得存在
|
||||
exists, taskType, goodsID, err := repository.TaskExists(tx, taskID)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if !exists {
|
||||
return "", ErrTaskNotFound
|
||||
}
|
||||
|
||||
// 3. 权限:**只有从没领过的才拒绝**。
|
||||
// 任务已取消、已重派给别人,都照样接受。
|
||||
claimed, err := repository.HasEverClaimed(tx, taskID, clientID)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if !claimed {
|
||||
return "", ErrNeverClaimed
|
||||
}
|
||||
|
||||
// 4. 业务写入
|
||||
resp, err := apply(tx, taskType, goodsID)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
respJSON, err := json.Marshal(resp)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("序列化响应失败: %w", err)
|
||||
}
|
||||
|
||||
// 5. 记下幂等键。和业务写入在同一个事务里——
|
||||
// 分开写的话,业务写成功但幂等没记上,客户端重试会被重复处理。
|
||||
if err := repository.SaveIdempotent(tx, idemKey, hash, string(respJSON)); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
// 6. 刷新客户端活动时间。三个接口都要做,
|
||||
// 只在 claim 里做的话,长任务期间会被误判成离线。
|
||||
if err := repository.TouchClient(tx, clientID); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
if err := tx.Commit(); err != nil {
|
||||
return "", fmt.Errorf("提交事务失败: %w", err)
|
||||
}
|
||||
return string(respJSON), nil
|
||||
}
|
||||
@@ -0,0 +1,374 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"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}); 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)
|
||||
}
|
||||
}
|
||||
|
||||
// ── 采集任务的结果要落到商品级 ─────────────────────────
|
||||
|
||||
func TestSubmitResult_采集结果写入商品级(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
now := model.NowISO()
|
||||
|
||||
if _, err := db.Exec(`
|
||||
INSERT INTO shopee_products (goods_id, title, pdd_goods_url,
|
||||
collect_status, created_at, updated_at)
|
||||
VALUES ('G-1','测试商品','https://x/1','collecting',?,?)`, now, now); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := db.Exec(`
|
||||
INSERT INTO tasks (task_id, task_type, status, assigned_client, goods_id,
|
||||
pdd_goods_url, created_at, updated_at)
|
||||
VALUES ('TASK-C','collect','assigned','client-001','G-1','https://x/1',?,?)`,
|
||||
now, now); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
claimTask(t, db, "TASK-C", "client-001")
|
||||
|
||||
body := `{"task_version":1,"attempt_id":"a-1","result_type":"collect","pdd_data":{"schema_version":1,"skus":[]}}`
|
||||
if _, err := SubmitResult(db, "TASK-C", "client-001", "key-c", []byte(body)); err != nil {
|
||||
t.Fatalf("提交采集结果失败: %v", err)
|
||||
}
|
||||
|
||||
var status, data string
|
||||
if err := db.QueryRow(
|
||||
`SELECT collect_status, COALESCE(pdd_data,'') FROM shopee_products WHERE goods_id='G-1'`,
|
||||
).Scan(&status, &data); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if status != "collected" {
|
||||
t.Errorf("采集状态应为 collected,实际 %s", status)
|
||||
}
|
||||
if data == "" {
|
||||
t.Error("pdd_data 应被写入商品级,实际为空")
|
||||
}
|
||||
}
|
||||
|
||||
// ── 提交失败 ───────────────────────────────────────────
|
||||
|
||||
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":"页面超时"}}`
|
||||
if _, err := SubmitFailure(db, "TASK-A", "client-001", "k", []byte(body)); err != nil {
|
||||
t.Fatalf("提交失败结果出错: %v", err)
|
||||
}
|
||||
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_采集失败写回商品级(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
now := model.NowISO()
|
||||
db.Exec(`INSERT INTO shopee_products (goods_id,title,collect_status,created_at,updated_at)
|
||||
VALUES ('G-1','测试商品','collecting',?,?)`, now, now)
|
||||
db.Exec(`INSERT INTO tasks (task_id,task_type,status,assigned_client,goods_id,
|
||||
pdd_goods_url,created_at,updated_at)
|
||||
VALUES ('TASK-C','collect','assigned','client-001','G-1','https://x/1',?,?)`, now, now)
|
||||
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)
|
||||
}
|
||||
|
||||
var status, errMsg string
|
||||
db.QueryRow(`SELECT collect_status, COALESCE(collect_error,'')
|
||||
FROM shopee_products WHERE goods_id='G-1'`).Scan(&status, &errMsg)
|
||||
if status != "failed" {
|
||||
t.Errorf("采集状态应为 failed,实际 %s", status)
|
||||
}
|
||||
if errMsg == "" {
|
||||
t.Error("失败原因应写入 collect_error,否则操作员看不到为什么采不到")
|
||||
}
|
||||
}
|
||||
|
||||
// ── 活动时间 ───────────────────────────────────────────
|
||||
|
||||
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 里刷新的话," +
|
||||
"客户端执行长任务期间会被误判成离线")
|
||||
}
|
||||
}
|
||||
@@ -36,6 +36,7 @@ db, _ := sql.Open("sqlite", dsn)
|
||||
| `busy_timeout(5000)` | 拿不到锁时最多等 5 秒,而不是立刻报错 |
|
||||
| `journal_mode(WAL)` | 读和写可以同时进行,不互相锁死 |
|
||||
| `foreign_keys(1)` | 打开外键约束(SQLite 默认是**关**的) |
|
||||
| `_txlock=immediate` | 事务一开始就拿写锁,见下 |
|
||||
|
||||
**为什么不能用 `db.Exec`:** Go 的 `database/sql` 是一个**连接池**。
|
||||
`db.Exec("PRAGMA busy_timeout=5000")` 只作用于当时拿到的那一条连接,
|
||||
@@ -46,6 +47,20 @@ db, _ := sql.Open("sqlite", dsn)
|
||||
这个坑在开发时不容易发现——单线程跑一切正常,一并发就炸。
|
||||
本项目的并发领取测试就是被它绊倒过一次。
|
||||
|
||||
**为什么必须加 `_txlock=immediate`:** Go 的 `db.Begin()` 默认发的是
|
||||
`BEGIN DEFERRED`——事务开始时**不拿写锁**,等第一次写才去拿。
|
||||
于是多个事务能同时开始、各自先读,然后同时想升级成写,互相卡死。
|
||||
这种情况 `busy_timeout` **救不了**,等下去也不会有结果。
|
||||
|
||||
实测(6 个并发事务,每个先读后写):
|
||||
|
||||
| DSN | 失败数 |
|
||||
|---|---|
|
||||
| 默认 deferred | **5 / 6** |
|
||||
| 加 `_txlock=immediate` | **0 / 6** |
|
||||
|
||||
加上之后事务一开始就排队拿锁,拿不到就按 `busy_timeout` 等,这才是要的行为。
|
||||
|
||||
`[建议]` 同时限制连接数:
|
||||
|
||||
```go
|
||||
@@ -322,7 +337,33 @@ CREATE INDEX idx_tasks_order ON tasks(order_no);
|
||||
|
||||
状态含义见 [01 需求](01-requirements.md) §6.2。
|
||||
|
||||
## 7. `clients` 客户端
|
||||
## 7. `task_claims` 领取历史
|
||||
|
||||
记录"哪台客户端领过哪个任务"。
|
||||
|
||||
```sql
|
||||
CREATE TABLE task_claims (
|
||||
task_id TEXT NOT NULL,
|
||||
client_id TEXT NOT NULL,
|
||||
claimed_at TEXT NOT NULL,
|
||||
PRIMARY KEY (task_id, client_id)
|
||||
);
|
||||
|
||||
CREATE INDEX idx_task_claims_client ON task_claims(client_id);
|
||||
```
|
||||
|
||||
**为什么需要这张表:** [04 Client 接口实现](04-client-api.md) §4.1 要求
|
||||
"只有**从未分配给该客户端**的任务才返回 403"。
|
||||
但 `tasks.assigned_client` 只记**当前**归属,任务一旦重派给别人,
|
||||
就查不出原来那台领过了——而契约又明确要求
|
||||
"**已重派仍要接受原客户端提交的结果**"。没有这张表,那条规则根本没法判断。
|
||||
|
||||
顺带得到一份审计记录:这个任务被哪几台客户端领过。
|
||||
|
||||
`[必须]` 领取成功时写入;提交结果时用它做权限判断。
|
||||
同一客户端重复领同一任务只更新时间,不报错。
|
||||
|
||||
## 8. `clients` 客户端
|
||||
|
||||
```sql
|
||||
CREATE TABLE clients (
|
||||
@@ -344,7 +385,7 @@ CREATE TABLE clients (
|
||||
- `[必须]` 注册发生在**领取任务时**,新序列号新增、已有的更新,
|
||||
**不设单独的注册或心跳接口**,见 [04](04-client-api.md) §3。
|
||||
|
||||
## 8. 数据关系总览
|
||||
## 9. 数据关系总览
|
||||
|
||||
```text
|
||||
shopee_products ──1:N──→ shopee_skus
|
||||
@@ -359,7 +400,7 @@ syb_orders ────────────────────┘
|
||||
└──创建──→ tasks ──分配──→ clients
|
||||
```
|
||||
|
||||
## 9. 与 Client 数据模型的关系
|
||||
## 10. 与 Client 数据模型的关系
|
||||
|
||||
Admin 和 Client **各有一个 SQLite,互不相通**,只通过接口交换数据。
|
||||
|
||||
|
||||
Reference in New Issue
Block a user