feat: add Chrome order backfill workflow (#316)

This commit is contained in:
QiuSW
2026-09-18 16:03:27 +08:00
parent 65d1f34865
commit 10be37498d
27 changed files with 876 additions and 20 deletions
@@ -91,6 +91,43 @@ func TestEveryRouteIsExplicitlyScoped(t *testing.T) {
}
}
func keyFor(e Endpoint) string { return e.Method + " " + e.Path }
func TestOrderBackfillRequiresExplicitWritebackAction(t *testing.T) {
db, s := fixture(t)
e := Endpoint{Method: "POST", Path: "/purchase-tasks/order-backfill", Module: "purchase_tasks", Capability: "writeback"}
router := gin.New()
called := 0
router.POST("/api/client/v1"+e.Path, Gate(db, e), func(c *gin.Context) { called++; c.JSON(200, gin.H{"data": "ok"}) })
send := func(token string) int {
req := httptest.NewRequest("POST", "https://example.test/api/client/v1"+e.Path, strings.NewReader(`{"requestId":"00000000-0000-4000-8000-000000000316","items":[]}`))
if token != "" {
req.Header.Set("Authorization", "Bearer "+token)
}
rr := httptest.NewRecorder()
router.ServeHTTP(rr, req)
return rr.Code
}
if got := send(""); got != 401 {
t.Fatalf("missing key status %d", got)
}
if (clientkey.Key{GrantsJSON: `[{"module":"purchase_tasks","write":true}]`}).Allows("purchase_tasks", "writeback") {
t.Fatal("generic write must not imply writeback")
}
_, writeOnly, err := s.Create(context.Background(), "purchase-only", []clientkey.Grant{{Module: "purchase_tasks", Actions: []string{"purchase"}}}, 1)
if err != nil {
t.Fatal(err)
}
if got := send(writeOnly); got != 403 {
t.Fatalf("generic write unexpectedly granted writeback: %d", got)
}
_, authorized, err := s.Create(context.Background(), "writeback", []clientkey.Grant{{Module: "purchase_tasks", Actions: []string{"writeback"}}}, 1)
if err != nil {
t.Fatal(err)
}
if got := send(authorized); got != 200 || called != 1 {
t.Fatalf("explicit writeback status=%d called=%d", got, called)
}
}
func TestRevocationTransportRedactionAndAudit(t *testing.T) {
db, s := fixture(t)
v, token, err := s.Create(context.Background(), "test", []clientkey.Grant{{Module: "pdd_products"}}, 1)
+1
View File
@@ -105,6 +105,7 @@ func Inventory() []Endpoint {
{"POST", "/purchase-tasks/batch", access.ModulePurchaseTasks, "purchase", buy.AdminBatchCreate},
{"POST", "/purchase-tasks/batch-retry", access.ModulePurchaseTasks, "purchase", buy.AdminBatchRetry},
{"POST", "/purchase-tasks/stock", access.ModulePurchaseTasks, "purchase", buy.AdminCreateStock},
{"POST", "/purchase-tasks/order-backfill", access.ModulePurchaseTasks, "writeback", buy.ClientBackfillOrders},
{"GET", "/ai-matching-settings", access.ModuleAIMatching, "read", func(c *gin.Context) {
db, err := pkg.GetOrm(c)
if err != nil {
+55 -11
View File
@@ -61,18 +61,43 @@ type OrderBackfillResponse struct {
Items []OrderBackfillResult `json:"items"`
}
type orderBackfillScope struct{ deviceID *uint64 }
func agentBackfillScope(deviceID uint64) orderBackfillScope {
return orderBackfillScope{deviceID: &deviceID}
}
func clientBackfillScope() orderBackfillScope { return orderBackfillScope{} }
func (scope orderBackfillScope) allows(task models.PurchaseTask) bool {
return scope.deviceID == nil || (task.DeviceID != nil && *task.DeviceID == *scope.deviceID)
}
func (scope orderBackfillScope) readable(db *gorm.DB, taskID uint64, saved *models.PurchaseTask) bool {
query := db.Where("id = ?", taskID)
if scope.deviceID != nil {
query = query.Where("device_id = ?", *scope.deviceID)
}
return query.First(saved).Error == nil
}
func (s *Service) BackfillOrders(ctx context.Context, req OrderBackfillRequest, token string) (OrderBackfillResponse, error) {
out := OrderBackfillResponse{RequestID: req.RequestID}
d, err := device.NewService(s.DB).Authenticate(ctx, token)
if err != nil {
return out, err
}
return s.backfillOrders(ctx, req, agentBackfillScope(d.ID))
}
func (s *Service) BackfillOrdersForClient(ctx context.Context, req OrderBackfillRequest) (OrderBackfillResponse, error) {
return s.backfillOrders(ctx, req, clientBackfillScope())
}
func (s *Service) backfillOrders(ctx context.Context, req OrderBackfillRequest, scope orderBackfillScope) (OrderBackfillResponse, error) {
out := OrderBackfillResponse{RequestID: req.RequestID}
if _, err := uuid.Parse(req.RequestID); err != nil || len(req.Items) == 0 || len(req.Items) > MaxOrderBackfillItems {
return out, fail(CodeInvalidRequest, "requestId 必须为 UUID,items 必须包含 1 到 50 条")
}
ids := make([]uint64, len(req.Items))
orders := make(map[uint64]string)
conflicts := make(map[uint64]bool)
orders, conflicts := make(map[uint64]string), make(map[uint64]bool)
for i, item := range req.Items {
id, err := purchasecontract.ParseAddressSuffix(item.AddressSuffix)
if err != nil {
@@ -90,7 +115,7 @@ func (s *Service) BackfillOrders(ctx context.Context, req OrderBackfillRequest,
if ids[i] == 0 {
r.Code = CodeBackfillSuffix
} else {
r = s.backfillOrder(ctx, d.ID, ids[i], req.RequestID, item, conflicts[ids[i]])
r = s.backfillOrder(ctx, scope, ids[i], req.RequestID, item, conflicts[ids[i]])
r.Index = i
}
out.Items[i] = r
@@ -98,7 +123,7 @@ func (s *Service) BackfillOrders(ctx context.Context, req OrderBackfillRequest,
return out, nil
}
func (s *Service) backfillOrder(ctx context.Context, deviceID, taskID uint64, requestID string, item OrderBackfillItem, batchConflict bool) OrderBackfillResult {
func (s *Service) backfillOrder(ctx context.Context, scope orderBackfillScope, taskID uint64, requestID string, item OrderBackfillItem, batchConflict bool) OrderBackfillResult {
r := OrderBackfillResult{TaskID: taskID, Result: "failed"}
var task models.PurchaseTask
// SQL errors must not print bound order numbers or the task's address snapshot.
@@ -107,7 +132,7 @@ func (s *Service) backfillOrder(ctx context.Context, deviceID, taskID uint64, re
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&task, taskID).Error; err != nil {
return purchaseNotFound(err)
}
if task.DeviceID == nil || *task.DeviceID != deviceID {
if !scope.allows(task) {
return fail(CodeBackfillDevice, "任务不属于当前设备")
}
if batchConflict {
@@ -133,18 +158,37 @@ func (s *Service) backfillOrder(ctx context.Context, deviceID, taskID uint64, re
if task.PDDOrderNo == nil || *task.PDDOrderNo != item.PDDOrderNo {
return fail(CodeStateConflict, "已创建订单缺少匹配订单号")
}
changed := false
filledSubmittedAt := false
if task.OrderSubmittedAt == nil && item.OrderSubmittedAt != nil {
submitted, parseErr := time.Parse(time.RFC3339Nano, *item.OrderSubmittedAt)
if parseErr != nil || submitted.IsZero() || submitted.Year() < 1000 || submitted.Year() > 9999 {
return fail(CodeOrderTimeInvalid, "下单时间必须为 RFC3339")
}
submitted = submitted.UTC()
task.OrderSubmittedAt = &submitted
changed = true
filledSubmittedAt = true
}
if item.PDDOrderAmountCent != nil {
if task.PDDOrderAmountCent == nil {
task.PDDOrderAmountCent = item.PDDOrderAmountCent
task.StatusVersion++
task.StatusChangedAt = s.Now()
if err := tx.Save(&task).Error; err != nil {
return err
}
changed = true
} else if *task.PDDOrderAmountCent != *item.PDDOrderAmountCent {
r.WarningCode, r.WarningMessage = CodeBackfillAmountConflict, "已有不同实付金额,本次未覆盖"
}
}
if changed {
if filledSubmittedAt {
marker := "backfill:page:" + uuid.NewSHA1(uuid.NameSpaceOID, []byte(requestID+":"+item.AddressSuffix)).String()
task.UnknownResolveRequestID = &marker
}
task.StatusVersion++
task.StatusChangedAt = s.Now()
if err := tx.Save(&task).Error; err != nil {
return err
}
}
r.Result, r.Code = "already_backfilled", "ALREADY_BACKFILLED"
return ensureOrderWriteback(tx, task)
}
@@ -208,7 +252,7 @@ func (s *Service) backfillOrder(ctx context.Context, deviceID, taskID uint64, re
readable := err == nil
if !readable {
saved = models.PurchaseTask{}
readable = db.Where("id = ? AND device_id = ?", taskID, deviceID).First(&saved).Error == nil
readable = scope.readable(db, taskID, &saved)
}
if readable {
r.Status, r.StatusVersion = saved.Status, saved.StatusVersion
@@ -1,6 +1,7 @@
package purchase
import (
"go-admin/common/clientprincipal"
"net/http"
"github.com/gin-gonic/gin"
@@ -23,3 +24,25 @@ func (h Handler) BackfillOrders(c *gin.Context) {
c.Header("Cache-Control", "no-store")
c.JSON(http.StatusOK, gin.H{"data": out})
}
func (h Handler) ClientBackfillOrders(c *gin.Context) {
if _, ok := clientprincipal.Get(c); !ok {
c.JSON(http.StatusUnauthorized, gin.H{"code": http.StatusUnauthorized, "message": "需要客户端身份"})
return
}
var req OrderBackfillRequest
if !decode(c, &req) {
return
}
s, ok := h.service(c)
if !ok {
return
}
out, err := s.BackfillOrdersForClient(c.Request.Context(), req)
if err != nil {
writeError(c, err)
return
}
c.Header("Cache-Control", "no-store")
c.JSON(http.StatusOK, gin.H{"data": out})
}
@@ -402,6 +402,55 @@ func TestOrderBackfillHTTPBoundary(t *testing.T) {
}
}
func TestClientOrderBackfillDoesNotRequireDeviceOwnership(t *testing.T) {
db := testDB(t)
f := seed(t, db, liveCaps(), true)
task := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
if err := db.Model(&task).Update("device_id", nil).Error; err != nil {
t.Fatal(err)
}
out, err := testService(db).BackfillOrdersForClient(context.Background(), OrderBackfillRequest{
RequestID: uuid.NewString(), Items: []OrderBackfillItem{backfillItem(task.ID, "CLIENT-ORDER")},
})
if err != nil || len(out.Items) != 1 || out.Items[0].Code != "BACKFILLED" {
t.Fatalf("client backfill: %+v %v", out, err)
}
}
func TestClientOrderBackfillFillsOnlyMissingSubmittedTime(t *testing.T) {
db := testDB(t)
f := seed(t, db, liveCaps(), true)
s := testService(db)
task := backfillTask(t, db, f, models.PurchaseTaskStatusOrderCreated)
order := "EXISTING-ORDER"
task.PDDOrderNo, task.OrderSubmittedAt = &order, nil
if err := db.Save(&task).Error; err != nil {
t.Fatal(err)
}
first := "2026-09-18T10:00:00+08:00"
item := backfillItem(task.ID, order)
item.OrderSubmittedAt = &first
out, err := s.BackfillOrdersForClient(context.Background(), OrderBackfillRequest{RequestID: uuid.NewString(), Items: []OrderBackfillItem{item}})
if err != nil || out.Items[0].Code != "ALREADY_BACKFILLED" || out.Items[0].TimeSource != "page" {
t.Fatalf("fill time: %+v %v", out, err)
}
want, _ := time.Parse(time.RFC3339, first)
got := loadBackfillTask(t, db, task.ID)
if got.OrderSubmittedAt == nil || !got.OrderSubmittedAt.Equal(want) {
t.Fatalf("missing time not filled: %+v", got.OrderSubmittedAt)
}
second := "2026-09-19T10:00:00+08:00"
item.OrderSubmittedAt = &second
out, err = s.BackfillOrdersForClient(context.Background(), OrderBackfillRequest{RequestID: uuid.NewString(), Items: []OrderBackfillItem{item}})
if err != nil || out.Items[0].TimeSource != "page" {
t.Fatalf("preserve source: %+v %v", out, err)
}
got = loadBackfillTask(t, db, task.ID)
if got.OrderSubmittedAt == nil || !got.OrderSubmittedAt.Equal(want) {
t.Fatalf("existing time overwritten: %+v", got.OrderSubmittedAt)
}
}
func TestOrderBackfillConcurrentSameOrderDifferentTasks(t *testing.T) {
db := testDB(t)
f := seed(t, db, liveCaps(), true)
@@ -223,7 +223,7 @@ func TestOrderWritebackEnqueueRollbackAndSameOrderReplay(t *testing.T) {
var dev models.AgentDevice
s.DB.First(&dev, *task.DeviceID)
// Directly exercise the already_backfilled branch without needing a raw token.
r := s.backfillOrder(context.Background(), dev.ID, task.ID, uuid.NewString(), backfillItem(task.ID, *task.PDDOrderNo), false)
r := s.backfillOrder(context.Background(), agentBackfillScope(dev.ID), task.ID, uuid.NewString(), backfillItem(task.ID, *task.PDDOrderNo), false)
if r.Result != "already_backfilled" {
t.Fatal(r.Code)
}