167 lines
5.8 KiB
Go
167 lines
5.8 KiB
Go
// 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"
|
|||
|
|
"net/http"
|
|||
|
|
|
|||
|
|
"github.com/gin-gonic/gin"
|
|||
|
|
)
|
|||
|
|
|
|||
|
|
// 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.POST("/tasks/claim", h.Claim)
|
|||
|
|
g.POST("/tasks/:task_id/result", h.SubmitResult)
|
|||
|
|
g.POST("/tasks/:task_id/failure", h.SubmitFailure)
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// ClaimRequest 是领取任务的请求体。
|
|||
|
|
// 字段对应 docs/client/04-admin-api-contract.md §5。
|
|||
|
|
type ClaimRequest struct {
|
|||
|
|
Client struct {
|
|||
|
|
// Name 是对 Client 契约的一处小扩展,需联合评审确认。
|
|||
|
|
// **要容忍它缺失**:没有就用 X-Client-Id 当显示名。
|
|||
|
|
Name string `json:"name"`
|
|||
|
|
} `json:"client"`
|
|||
|
|
SupportedTypes []string `json:"supported_types"`
|
|||
|
|
Device struct {
|
|||
|
|
Address string `json:"address"`
|
|||
|
|
Platform string `json:"platform"`
|
|||
|
|
PddPackage string `json:"pdd_package"`
|
|||
|
|
} `json:"device"`
|
|||
|
|
Capabilities struct {
|
|||
|
|
PurchaseMode string `json:"purchase_mode"`
|
|||
|
|
SchemaVersions []int `json:"schema_versions"`
|
|||
|
|
} `json:"capabilities"`
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// 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 ClaimRequest
|
|||
|
|
if err := c.ShouldBindJSON(&req); err != nil {
|
|||
|
|
apiError(c, http.StatusBadRequest, "INVALID_BODY",
|
|||
|
|
"请求体不是合法 JSON", false)
|
|||
|
|
return
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// TODO(骨架): 1. service.RegisterClient(h.db, clientID, req)
|
|||
|
|
// 新 client_id 就新增,已有就更新 device/capabilities/last_seen_at。
|
|||
|
|
// name 若已被人工改过,不要覆盖。
|
|||
|
|
|
|||
|
|
// TODO(骨架): 2. service.ClaimNextTask(h.db, clientID)
|
|||
|
|
// 只返回分配给这个客户端的任务,一次一条。
|
|||
|
|
// UPDATE ... WHERE task_id = ? AND status = 'assigned'
|
|||
|
|
// 检查影响行数,为 0 说明被抢先了,取下一条。
|
|||
|
|
|
|||
|
|
// 新客户端第一次来必然拿不到任务(还没人给它分配),返回 204 是正常的。
|
|||
|
|
c.Status(http.StatusNoContent)
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// SubmitResult 接收成功结果。
|
|||
|
|
//
|
|||
|
|
// **最容易写错的一条**(docs/admin/04-client-api.md §4.1):
|
|||
|
|
// - 不得因为任务已取消而拒绝;
|
|||
|
|
// - 不得因为任务已重派给别的客户端而拒绝;
|
|||
|
|
// - 必须能接受同一任务来自多个客户端的多份结果;
|
|||
|
|
// - 只有**从未分配给该客户端**的任务才返回 403。
|
|||
|
|
//
|
|||
|
|
// 原因:Client 中途不查任务状态,所以它必然会提交一些
|
|||
|
|
// "Admin 这边已经不要了"的结果。而它可能真的已经下单了,
|
|||
|
|
// 这些数据必须能交上来留痕。
|
|||
|
|
//
|
|||
|
|
// accepted: true 的意思是"我收到并存下了",不代表任务还算数。
|
|||
|
|
func (h *Handler) SubmitResult(c *gin.Context) {
|
|||
|
|
taskID := c.Param("task_id")
|
|||
|
|
idemKey := c.GetHeader("Idempotency-Key")
|
|||
|
|
if idemKey == "" {
|
|||
|
|
apiError(c, http.StatusBadRequest, "MISSING_IDEMPOTENCY_KEY",
|
|||
|
|
"缺少 Idempotency-Key 请求头", false)
|
|||
|
|
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,不刷新会被误判成离线)。
|
|||
|
|
|
|||
|
|
_ = taskID
|
|||
|
|
apiError(c, http.StatusNotImplemented, "NOT_IMPLEMENTED",
|
|||
|
|
"提交结果接口尚未实现", false)
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// 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")
|
|||
|
|
|
|||
|
|
// TODO(骨架): 同 SubmitResult 的幂等处理;
|
|||
|
|
// 采集任务失败时同步把 shopee_products.collect_status 置 failed,
|
|||
|
|
// 错误信息写 collect_error,界面上要看得见。
|
|||
|
|
|
|||
|
|
_ = taskID
|
|||
|
|
apiError(c, http.StatusNotImplemented, "NOT_IMPLEMENTED",
|
|||
|
|
"提交失败接口尚未实现", false)
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// 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{},
|
|||
|
|
},
|
|||
|
|
})
|
|||
|
|
}
|