feat: add safe agent purchase retry (#95)

This commit is contained in:
QiuSW
2026-08-26 14:37:01 +08:00
parent 9860b2aa63
commit 5f2dd05ebd
13 changed files with 352 additions and 19 deletions
+10 -3
View File
@@ -37,6 +37,8 @@ type AgentPurchaseItem struct {
OrderSubmittedAt *time.Time `json:"orderSubmittedAt,omitempty"`
ErrorCode *string `json:"errorCode,omitempty"`
ErrorMessage *string `json:"errorMessage,omitempty"`
Retryable bool `json:"retryable"`
RetryDisabledReason string `json:"retryDisabledReason,omitempty"`
CreatedAt time.Time `json:"createdAt"`
}
@@ -82,7 +84,7 @@ func (s *Service) AgentHistory(ctx context.Context, req AgentHistoryRequest, tok
}
items := make([]AgentPurchaseItem, 0, len(tasks))
for _, task := range tasks {
items = append(items, agentPurchaseItem(task))
items = append(items, agentPurchaseItem(task, s.retryQueryEligibility(ctx, task, true)))
}
return AgentPurchaseList{Items: items, Total: total, Page: req.Page, PageSize: req.PageSize}, nil
}
@@ -102,10 +104,14 @@ func (s *Service) AgentHistoryDetail(ctx context.Context, taskID uint64, token s
if err != nil {
return AgentPurchaseDetail{}, internal(err)
}
return AgentPurchaseDetail{Task: agentPurchaseItem(task)}, nil
return AgentPurchaseDetail{Task: agentPurchaseItem(task, s.retryQueryEligibility(ctx, task, true))}, nil
}
func agentPurchaseItem(task models.PurchaseTask) AgentPurchaseItem {
func agentPurchaseItem(task models.PurchaseTask, retry retryDecision) AgentPurchaseItem {
retryDisabledReason := ""
if task.Status == models.PurchaseTaskStatusFailed {
retryDisabledReason = retry.Reason
}
return AgentPurchaseItem{
TaskID: task.ID, Status: task.Status, ShopeeOrderNo: task.ShopeeOrderNoSnapshot,
PDDGoodsID: task.PDDGoodsIDSnapshot, PDDTitle: task.PDDTitleSnapshot,
@@ -113,6 +119,7 @@ func agentPurchaseItem(task models.PurchaseTask) AgentPurchaseItem {
Quantity: task.Quantity, ActualUnitPriceCent: task.ActualUnitPriceCent, Currency: task.Currency,
PDDOrderNo: task.PDDOrderNo, OrderSubmittedAt: task.OrderSubmittedAt,
ErrorCode: task.ErrorCode, ErrorMessage: task.ErrorMessage, CreatedAt: task.CreatedAt,
Retryable: retry.Allowed, RetryDisabledReason: retryDisabledReason,
}
}
@@ -0,0 +1,79 @@
package purchase
import (
"context"
"testing"
"go-admin/app/goauto/device"
"go-admin/app/goauto/models"
"github.com/google/uuid"
"gorm.io/gorm"
)
func TestAgentRetryCreatesOneFixedDeviceTaskAndReplays(t *testing.T) {
db := testDB(t)
f := seed(t, db, liveCaps(), true)
setCollectedPDDPrice(t, db, f.pdd.ID)
service := testService(db)
failed := failedLiveTask(t, db, service, f)
history, err := service.AgentHistoryDetail(context.Background(), failed.ID, f.token)
if err != nil || !history.Task.Retryable || history.Task.RetryDisabledReason != "" {
t.Fatalf("safe failure was not exposed as retryable: %+v err=%v", history, err)
}
request := AgentRetryRequest{RequestID: uuid.NewString()}
first, err := service.AgentRetry(context.Background(), failed.ID, request, f.token)
if err != nil || first.TaskID == 0 || first.TaskID == failed.ID || first.SourceTaskID != failed.ID || first.Replayed {
t.Fatalf("unexpected first retry: %+v err=%v", first, err)
}
var oldTask, newTask models.PurchaseTask
if err = db.First(&oldTask, failed.ID).Error; err != nil {
t.Fatal(err)
}
if err = db.First(&newTask, first.TaskID).Error; err != nil {
t.Fatal(err)
}
if oldTask.Status != models.PurchaseTaskStatusFailed || newTask.Status != models.PurchaseTaskStatusPending || newTask.DeviceID == nil || *newTask.DeviceID != f.device.ID || newTask.AddressSuffix == oldTask.AddressSuffix {
t.Fatalf("retry did not preserve old task or fix new task to device: old=%+v new=%+v", oldTask, newTask)
}
replay, err := service.AgentRetry(context.Background(), failed.ID, request, f.token)
if err != nil || !replay.Replayed || replay.TaskID != first.TaskID {
t.Fatalf("request id was not idempotent: first=%+v replay=%+v err=%v", first, replay, err)
}
}
func TestAgentRetryRejectsAnotherDeviceAndUnsafeBoundary(t *testing.T) {
db := testDB(t)
f := seed(t, db, liveCaps(), true)
setCollectedPDDPrice(t, db, f.pdd.ID)
service := testService(db)
failed := failedLiveTask(t, db, service, f)
otherToken := uuid.NewString()
registration := device.NewService(db)
registration.GenerateToken = func() (string, error) { return otherToken, nil }
if _, err := registration.Register(context.Background(), device.RegisterRequest{
RequestID: uuid.NewString(), InstallID: uuid.NewString(), Name: "other", Manufacturer: "Samsung",
Model: "Test", AndroidVersion: "14", AgentVersion: "1", PDDVersion: "7",
}, ""); err != nil {
t.Fatal(err)
}
if _, err := service.AgentRetry(context.Background(), failed.ID, AgentRetryRequest{RequestID: uuid.NewString()}, otherToken); code(err) != CodeTaskNotFound {
t.Fatalf("other device could retry task: %v", err)
}
irreversible := service.Now()
if err := db.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", failed.ID).Update("irreversible_at", irreversible).Error; err != nil {
t.Fatal(err)
}
detail, err := service.AgentHistoryDetail(context.Background(), failed.ID, f.token)
if err != nil || detail.Task.Retryable || detail.Task.RetryDisabledReason == "" {
t.Fatalf("unsafe task eligibility mismatch: %+v err=%v", detail, err)
}
if _, err := service.AgentRetry(context.Background(), failed.ID, AgentRetryRequest{RequestID: uuid.NewString()}, f.token); code(err) != CodeRetryUnsafe {
t.Fatalf("unsafe boundary was retryable: %v", err)
}
}
+21
View File
@@ -218,6 +218,27 @@ func (h Handler) AgentHistoryDetail(c *gin.Context) {
c.Header("Cache-Control", "no-store")
c.JSON(http.StatusOK, gin.H{"data": result})
}
func (h Handler) AgentRetry(c *gin.Context) {
id, ok := pathID(c)
if !ok {
return
}
var req AgentRetryRequest
if !decode(c, &req) {
return
}
service, serviceOK := h.service(c)
if !serviceOK {
return
}
result, err := service.AgentRetry(c.Request.Context(), id, req, bearer(c.GetHeader("Authorization")))
if err != nil {
writeError(c, err)
return
}
c.Header("Cache-Control", "no-store")
c.JSON(http.StatusOK, gin.H{"data": result})
}
func (h Handler) Claim(c *gin.Context) { h.action(c, (*Service).Claim) }
func (h Handler) Start(c *gin.Context) { h.action(c, (*Service).Start) }
func (h Handler) OrderSubmitStarted(c *gin.Context) { h.action(c, (*Service).MarkOrderSubmitStarted) }
+51
View File
@@ -6,6 +6,7 @@ import (
"fmt"
"strings"
"go-admin/app/goauto/device"
"go-admin/app/goauto/models"
"go-admin/app/goauto/purchasecontract"
@@ -37,6 +38,18 @@ type BatchRetryResponse struct {
FailedCount int `json:"failedCount"`
}
type AgentRetryRequest struct {
RequestID string `json:"requestId"`
}
type AgentRetryResponse struct {
SourceTaskID uint64 `json:"sourceTaskId"`
SourceTaskNo string `json:"sourceTaskNo"`
TaskID uint64 `json:"taskId"`
TaskNo string `json:"taskNo"`
Replayed bool `json:"replayed"`
}
type retryDecision struct {
Allowed bool
ReasonCode string
@@ -127,6 +140,44 @@ func (s *Service) BatchRetry(ctx context.Context, req BatchRetryRequest) (BatchR
return response, nil
}
// AgentRetry creates one safe replacement task for the authenticated device.
// It delegates creation and idempotency to BatchRetry so Admin and Agent use
// exactly the same archive, device and irreversible-boundary checks.
func (s *Service) AgentRetry(ctx context.Context, taskID uint64, req AgentRetryRequest, token string) (AgentRetryResponse, error) {
deviceRecord, err := device.NewService(s.DB).Authenticate(ctx, token)
if err != nil {
return AgentRetryResponse{}, err
}
if taskID == 0 {
return AgentRetryResponse{}, fail(CodeTaskNotFound, "采购任务不存在")
}
var source models.PurchaseTask
queryErr := s.DB.WithContext(ctx).
Where("id = ? AND device_id = ? AND created_at >= ?", taskID, deviceRecord.ID, s.Now().AddDate(0, 0, -agentPurchaseHistoryDays)).
First(&source).Error
if errors.Is(queryErr, gorm.ErrRecordNotFound) {
return AgentRetryResponse{}, fail(CodeTaskNotFound, "采购任务不存在")
}
if queryErr != nil {
return AgentRetryResponse{}, internal(queryErr)
}
result, err := s.BatchRetry(ctx, BatchRetryRequest{RequestID: req.RequestID, TaskIDs: []uint64{taskID}})
if err != nil {
return AgentRetryResponse{}, err
}
if len(result.Items) != 1 || !result.Items[0].Created || result.Items[0].TaskID == nil {
if len(result.Items) == 1 && result.Items[0].ReasonCode != "" {
return AgentRetryResponse{}, fail(result.Items[0].ReasonCode, result.Items[0].Reason)
}
return AgentRetryResponse{}, fail(CodeInternal, "服务端处理失败")
}
item := result.Items[0]
return AgentRetryResponse{
SourceTaskID: item.SourceTaskID, SourceTaskNo: item.SourceTaskNo,
TaskID: *item.TaskID, TaskNo: item.TaskNo, Replayed: item.Replayed,
}, nil
}
func validateRetryTaskIDs(raw []uint64) ([]uint64, error) {
if len(raw) == 0 || len(raw) > maxBatchRetryItems {
return nil, fail(CodeInvalidRequest, "taskIds 必须包含 1 至 100 条当前页任务")
+1
View File
@@ -19,6 +19,7 @@ func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) {
agent.GET("", h.AgentHistory)
agent.GET("/next", h.Next)
agent.GET("/:taskId", h.AgentHistoryDetail)
agent.POST("/:taskId/retry", h.AgentRetry)
agent.POST("/:taskId/claim", h.Claim)
agent.POST("/:taskId/start", h.Start)
agent.POST("/:taskId/order-submit-started", h.OrderSubmitStarted)