Files
cmautobuy/admin/handler/api/client_api.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

283 lines
9.6 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 api 是给 Client 用的 JSON 接口。
//
// **权威契约在 Client 那边**:docs/client/04-admin-api-contract.md。
// 本包只是实现,两边说法不一致时以 Client 契约为准,要改先走工单。
//
// 一共只有三个接口,**不得新增"让 Client 查询状态"类接口**,
// 也不得加回租约和心跳,理由见 docs/admin/04-client-api.md §1。
package api
import (
"database/sql"
"encoding/json"
"errors"
"io"
"log"
"net/http"
"github.com/gin-gonic/gin"
"cmautobuy/admin/model"
"cmautobuy/admin/service"
)
// Handler 持有接口共用的依赖。
type Handler struct {
db *sql.DB
}
// Register 挂上三个接口。
func Register(r *gin.Engine, db *sql.DB) {
h := &Handler{db: db}
g := r.Group("/api/v1/client")
{
// 登记:只登记客户端,**不碰任务**。设置页点"保存"调它。
g.PUT("/registration", h.Register)
// 领取:可能改变任务状态,只有"获取任务"流程能调。
g.POST("/tasks/claim", h.Claim)
g.POST("/tasks/:task_id/result", h.SubmitResult)
g.POST("/tasks/:task_id/failure", h.SubmitFailure)
}
}
// Register 处理设置页发起的显式登记。
//
// `[必须]` 它**不读取、不领取、不修改任何任务**,也不返回任务。
// 设置页点"保存"不该顺带把一个任务领走——领走了 Admin 就标成 claimed,
// 而保存动作没有义务去可靠保存那个任务,任务就丢了。
//
// 相同 X-Client-Id 重复调用是幂等的,不会产生重复记录。
func (h *Handler) Register(c *gin.Context) {
clientID := c.GetHeader("X-Client-Id")
if clientID == "" {
apiError(c, http.StatusBadRequest, "MISSING_CLIENT_ID",
"缺少 X-Client-Id 请求头", false)
return
}
var req service.ClientProfileRequest
if err := c.ShouldBindJSON(&req); err != nil {
apiError(c, http.StatusBadRequest, "INVALID_BODY",
"请求体不是合法 JSON", false)
return
}
registeredAt, err := service.RegisterClientProfile(h.db, clientID, req)
switch {
case err == nil:
log.Printf("client_registered client_id=%s", clientID)
c.JSON(http.StatusOK, gin.H{
"registered": true,
"client_id": clientID,
"registered_at": registeredAt,
})
case errors.Is(err, service.ErrInvalidProfile):
// 422:JSON 是合法的,但字段内容不符合规则
apiError(c, http.StatusUnprocessableEntity, "INVALID_CLIENT_PROFILE",
err.Error(), false)
default:
log.Printf("client_register_failed client_id=%s err=%v", clientID, err)
apiError(c, http.StatusInternalServerError, "CLIENT_REGISTER_FAILED",
"登记客户端失败,请稍后重试", true)
}
}
// Claim 领取一个任务。
//
// 处理顺序(见 docs/admin/04-client-api.md §2):
// 1. 先做客户端注册/更新——**注册就在这里做,没有单独的注册接口**;
// 2. 查 assigned_client = X-Client-Id 且 status = 'assigned' 的任务,取一条;
// 3. 没有返回 204 No Content(不是 200 加空对象);
// 4. 有就用条件更新把状态改成 claimed,检查影响行数防并发。
//
// 响应**不含租约**,也不含 Admin 侧状态。
func (h *Handler) Claim(c *gin.Context) {
clientID := c.GetHeader("X-Client-Id")
if clientID == "" {
apiError(c, http.StatusBadRequest, "MISSING_CLIENT_ID",
"缺少 X-Client-Id 请求头", false)
return
}
// 和登记接口共用同一个结构和校验,避免两处漂移
var req service.ClientProfileRequest
if err := c.ShouldBindJSON(&req); err != nil {
apiError(c, http.StatusBadRequest, "INVALID_BODY",
"请求体不是合法 JSON", false)
return
}
if err := req.Validate(); err != nil {
apiError(c, http.StatusUnprocessableEntity, "INVALID_CLIENT_PROFILE",
err.Error(), false)
return
}
// 1. 隐式登记。explicit=false —— 后台调用**不更新名称**,
// 否则操作员在 Admin 改的名字会被反复冲掉。
if err := service.RegisterClient(h.db, req.ToClient(clientID), false); err != nil {
log.Printf("client_register_failed client_id=%s err=%v", clientID, err)
apiError(c, http.StatusInternalServerError, "CLIENT_REGISTER_FAILED",
"登记客户端失败", true)
return
}
// 2. 领取一个任务,只会拿到分配给这个客户端的
task, err := service.ClaimNextTask(h.db, clientID, req.SupportedTypes)
if err != nil {
log.Printf("task_claim_failed client_id=%s err=%v", clientID, err)
apiError(c, http.StatusInternalServerError, "TASK_CLAIM_FAILED",
"领取任务失败", true)
return
}
// 3. 没有可领的任务 -> 204 No Content,**不是 200 加空对象**。
// 新客户端第一次来必然走到这里(还没人给它分配),是正常的。
if task == nil {
c.Status(http.StatusNoContent)
return
}
log.Printf("task_claimed client_id=%s task_id=%s type=%s",
clientID, task.TaskID, task.TaskType)
c.JSON(http.StatusOK, gin.H{"task": taskPayload(task)})
}
// taskPayload 把任务转成契约里的响应结构。
//
// **不含租约,也不含 Admin 侧状态**——Client 不关心这些,
// 见 docs/client/04-admin-api-contract.md §5。
func taskPayload(t *model.Task) gin.H {
payload := gin.H{
"goods_url": t.PddGoodsURL,
}
if t.PddGoodsID != "" {
payload["goods_id"] = t.PddGoodsID
}
if t.PddOptions != "" {
// pdd_options 存的是 JSON 字符串,要还原成对象再嵌进去,
// 否则 Client 拿到的是一个字符串而不是 {"color":...}
var options map[string]any
if err := json.Unmarshal([]byte(t.PddOptions), &options); err == nil {
payload["options"] = options
}
}
// 采购任务必须带数量和价格上限,这是 Client 的价格保护
if t.Quantity > 0 {
payload["quantity"] = t.Quantity
}
if t.MaxPriceCent > 0 {
payload["max_price_cent"] = t.MaxPriceCent
}
return gin.H{
"id": t.TaskID,
"type": t.TaskType,
"version": t.Version,
"priority": t.Priority,
"payload": payload,
"created_at": t.CreatedAt,
"updated_at": t.UpdatedAt,
}
}
// SubmitResult 接收成功结果。
//
// **最容易写错的一条**(docs/admin/04-client-api.md §4.1):
// - 不得因为任务已取消而拒绝;
// - 不得因为任务已重派给别的客户端而拒绝;
// - 必须能接受同一任务来自多个客户端的多份结果;
// - 只有**从未领过**这个任务的客户端才返回 403。
//
// 原因见 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",
"缺少 Idempotency-Key 请求头", false)
return
}
// 读**原始字节**而不是解析后再序列化——幂等要比对请求内容的哈希,
// 重新序列化会因为字段顺序、空格不同而算出不一样的哈希。
rawBody, err := io.ReadAll(c.Request.Body)
if err != nil {
apiError(c, http.StatusBadRequest, "INVALID_BODY", "读取请求体失败", false)
return
}
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))
case errors.Is(err, service.ErrTaskNotFound):
apiError(c, http.StatusNotFound, "TASK_NOT_FOUND",
"任务不存在", false)
case errors.Is(err, service.ErrNeverClaimed):
// 注意:只有**从没领过**才会走到这里。
// 任务已取消、已重派,都不会拒绝。
apiError(c, http.StatusForbidden, "TASK_NOT_ASSIGNED",
"该任务从未分配给这个客户端", false)
case errors.Is(err, service.ErrCollectMismatch):
// 采回来的商品不是请求的那个(链接跳转、采错商品)。
// 内容本身是合法 JSON,只是业务上对不上,所以是 422 不是 400。
apiError(c, http.StatusUnprocessableEntity, "COLLECT_GOODS_MISMATCH",
err.Error(), 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 按契约格式返回错误。
//
// 接口出错返回 JSON,页面出错渲染错误页,两者不要混。
// 格式见 docs/client/04-admin-api-contract.md §3。
func apiError(c *gin.Context, status int, code, message string, retryable bool) {
c.JSON(status, gin.H{
"error": gin.H{
"code": code,
"message": message,
"retryable": retryable,
"request_id": c.GetHeader("X-Request-Id"),
"details": gin.H{},
},
})
}