feat(purchase): add agent order backfill endpoint (#241)
This commit is contained in:
@@ -985,3 +985,76 @@ Agent 携带既有 Token(可已失效)及恢复码重新调用注册接口
|
||||
`orderCount` 是已验证页的原始列表读取数量;`detailCount`、`created`、`updated` 为已提交明细及其新增/覆盖数量;`daysProcessed` 是完整通过的日期数,不是已尝试日期数。失败日期/页码/阶段写入现有脱敏限长 errorMessage。部分成功不刷新店铺的完整同步统计。
|
||||
|
||||
Web 唯一展示位置为“采集采购 → SYB 同步记录”:列表状态、状态筛选及详情支持部分成功,详情保留已保存数量、错误原因与重新同步补齐提示。定时任务日志只表示异步任务受理,不等于最终业务同步成功。
|
||||
|
||||
## Agent 采购订单批量回填(#241)
|
||||
|
||||
本节为 #241 服务端实现契约,2026-09-08 按用户授权直接更新本地镜像;线上 Wiki 与其他长期文档由审核阶段同步。本节不表示已经部署或完成真机验收。
|
||||
|
||||
`POST /api/agent/v1/purchase-tasks/order-backfill`
|
||||
|
||||
使用 `Authorization: Bearer <Device Token>`,沿用 `RequireAgentHTTPS`、`GOAUTO_ALLOW_INSECURE_AGENT_HTTP` 与既有可信转发协议策略。无需 Admin JWT、claim、start 或 attempt。设备号只从认证读取,请求不得指定 deviceId、地址全文、收件人、手机号或原始控件树;未知 JSON 字段拒绝。此接口只记录已观察到的订单事实,不执行设备动作、创建订单或付款。
|
||||
|
||||
请求示例(页面时间先按 Asia/Shanghai 理解,再以带时区 RFC3339/RFC3339Nano 发送):
|
||||
|
||||
```json
|
||||
{
|
||||
"requestId": "5826cdda-dcd6-442e-90c3-9b75ba6fb8d8",
|
||||
"items": [
|
||||
{"addressSuffix": "_cg7", "pddOrderNo": "EXAMPLE-ORDER-7", "orderSubmittedAt": "2026-09-08T20:30:00+08:00"},
|
||||
{"addressSuffix": "_cg72", "pddOrderNo": "EXAMPLE-ORDER-72"}
|
||||
]
|
||||
}
|
||||
```
|
||||
|
||||
- requestId 必须为 UUID;items 为 1~50 条,保持输入顺序;请求体上限沿用 1 MiB。
|
||||
- addressSuffix 只接受 `AddressSuffix(id)` 生成的完整字符串。任务号为非零 uint64,拒绝前导零、正负号、空格、尾随文本、多个后缀与溢出;`_cg7` 与 `_cg72` 分别定位任务 7 和 72。
|
||||
- pddOrderNo 必填,最多 100 个 Unicode 字符,不接受首尾空白、换行或制表符,不自动裁剪后覆盖旧值。
|
||||
- orderSubmittedAt 缺失或 null 时回落任务 irreversible_at;空字符串、无时区文本及非法时间是条目错误,不触发回落。存储统一 UTC。页面值与 irreversible_at 都缺失时该条失败。
|
||||
|
||||
有效批次返回 HTTP 200,包括全部条目失败的批次;每条独立事务,失败不撤销其他条目已提交的数据。响应包裹为 `data`,并设 `Cache-Control: no-store`:
|
||||
|
||||
```json
|
||||
{
|
||||
"data": {
|
||||
"requestId": "5826cdda-dcd6-442e-90c3-9b75ba6fb8d8",
|
||||
"items": [
|
||||
{"index": 0, "taskId": 7, "result": "backfilled", "code": "BACKFILLED", "status": "order_created", "statusVersion": 5, "pddOrderNo": "EXAMPLE-ORDER-7", "orderSubmittedAt": "2026-09-08T12:30:00Z", "timeSource": "page", "retryable": false},
|
||||
{"index": 1, "taskId": 72, "result": "backfilled", "code": "BACKFILLED", "status": "order_created", "statusVersion": 4, "pddOrderNo": "EXAMPLE-ORDER-72", "orderSubmittedAt": "2026-09-08T12:31:00Z", "timeSource": "irreversible_at", "retryable": false}
|
||||
]
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
index 从 0 开始;后缀无法解析时不返回 taskId。result 为 `backfilled`、`already_backfilled`、`conflict` 或 `failed`。已认证设备所属任务可返回提交后的状态、版本、已保存订单号及时间;拒绝条目尽可能返回当前已提交事实。跨设备任务和不存在任务不返回这些业务字段,事务回滚后的内存值绝不作为最终事实返回。
|
||||
|
||||
timeSource 说明已保存时间的来源:`page` 为页面值,`irreversible_at` 为估算回落,`existing_unknown` 为原先已创建的历史订单且没有可证明的来源。没有已保存时间时省略 timeSource。客户端必须保留估算标记,不得把回落值或 unknown 宣称为页面真实时间。重复回填不会用新页面时间自动校正旧时间。
|
||||
|
||||
| 条目 code | result | 含义 |
|
||||
|---|---|---|
|
||||
| `BACKFILLED` | backfilled | 本条完成回填 |
|
||||
| `ALREADY_BACKFILLED` | already_backfilled | 正式任务已为 order_created 且订单号相同,无写入 |
|
||||
| `PURCHASE_BACKFILL_SUFFIX_INVALID` | failed | 非法、非规范、零、溢出或歧义后缀 |
|
||||
| `PURCHASE_TASK_NOT_FOUND` | failed | 任务不存在 |
|
||||
| `PURCHASE_BACKFILL_DEVICE_MISMATCH` | failed | 未绑定设备或不属于认证设备 |
|
||||
| `PURCHASE_STATE_CONFLICT` | failed | 非正式采购,或状态不允许回填 |
|
||||
| `PURCHASE_INVALID_REQUEST` | failed | 订单号非法 |
|
||||
| `PURCHASE_ORDER_TIME_INVALID` | failed | 提供的页面时间无效 |
|
||||
| `PURCHASE_ORDER_TIME_MISSING` | failed | 页面时间与 irreversible_at 均无有效值 |
|
||||
| `PURCHASE_BACKFILL_ORDER_CONFLICT` | conflict | 任务已有不同订单号 |
|
||||
| `PURCHASE_BACKFILL_BATCH_CONFLICT` | conflict | 同批同任务出现多个不同订单号,该任务所有条目均拒绝 |
|
||||
| `PURCHASE_BACKFILL_ORDER_ALREADY_USED` | conflict | 同一订单号已对应其他任务 |
|
||||
| `INTERNAL_ERROR` | failed | 数据库失败、死锁等,retryable=true,可安全重放 |
|
||||
|
||||
批级 JSON/UUID/数量错误为 HTTP 422 `PURCHASE_INVALID_REQUEST`;鉴权、停用设备和 HTTPS 限制复用既有错误(401 `DEVICE_TOKEN_INVALID`、403 `DEVICE_DISABLED`、426 `HTTPS_REQUIRED`)。批级失败使用既有 `{code,message,retryable}` 包裹,未开始条目写入。
|
||||
|
||||
### 状态、幂等与并发
|
||||
|
||||
新服务在事务中锁定任务并检查来源状态,仅允许当前设备的 `live + order_result_unknown` 首次写入;`live + order_created` 只在订单号相同时返回已回填。其他状态(包括 running、failed、cancelled 与演练)均拒绝。复用 SetStatus 同步占用字段,同一事务递增 statusVersion、设置 statusChangedAt、清空主任务当前错误及租约;原始 attempt、规则快照、支付与物流、SYB 回写字段不变。
|
||||
|
||||
requestId 沿用 UUID 约定,不增加批次表或全局幂等缓存。既有 unknown_resolve_request_id 槽存储 `backfill:<page|irreversible_at>:<由 requestId 和后缀派生的 UUID>`(最多 61 字符),用于任务级关联和保留本功能时间来源。相同任务和订单号即使更换 requestId 也无写入;同 requestId 改内容仍重新执行设备、状态和订单冲突检查,不凭 requestId 直接放行。批内不同任务提交同一订单号时,先成功提交者占用,其余条目返回订单已被使用;已存在的历史重复订单号不自动修复。
|
||||
|
||||
订单号尚无唯一索引,本实现不迁移数据库。共享模型保存钩子复用 `purchase_rule_setting.id=1` 行作短事务互斥锁,再以锁定读检查订单号归属,覆盖回填、原人工解除和旧结果提交路径;单例缺失时拒绝写入,数据库死锁时回滚失败事务。原 ResolveUnknown 与 Admin 鉴权代码保持不变。禁止通过跳过模型钩子的直接 SQL 写入宣称具备此保证。
|
||||
|
||||
新回填路径禁用包含绑定参数的 SQL 日志,不记录请求正文、订单号、地址或原始树;任务号和认证设备号沿既有任务关联不可变 ruleSnapshot。此接口未新增日志载荷或任务/attempt。
|
||||
|
||||
本地测试覆盖事务回滚、并发服务调用和旧写入路径,使用 SQLite;MySQL 8.4 多连接/多进程的实际行锁、生产数据和真机端到端回填尚待环境验收,不以单元测试替代。
|
||||
|
||||
@@ -202,7 +202,15 @@ func (task *PurchaseTask) BeforeCreate(_ *gorm.DB) error {
|
||||
return task.syncPurchaseGuardSlots()
|
||||
}
|
||||
|
||||
func (task *PurchaseTask) BeforeSave(_ *gorm.DB) error { return task.syncPurchaseGuardSlots() }
|
||||
func (task *PurchaseTask) BeforeSave(tx *gorm.DB) error {
|
||||
if err := task.syncPurchaseGuardSlots(); err != nil {
|
||||
return err
|
||||
}
|
||||
if task.PDDOrderNo != nil && *task.PDDOrderNo != "" {
|
||||
return CheckPurchaseOrderNumber(tx, task.ID, *task.PDDOrderNo)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (task *PurchaseTask) SetStatus(status string) error {
|
||||
task.Status = status
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
package models
|
||||
|
||||
import (
|
||||
"errors"
|
||||
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
var ErrPurchaseOrderNumberUsed = errors.New("purchase order number belongs to another task")
|
||||
|
||||
// CheckPurchaseOrderNumber must run inside the caller's write transaction.
|
||||
// The existing singleton setting row serializes order assignments across
|
||||
// processes, including an absent order number, without relying on gap locks or
|
||||
// a new schema constraint. Locking reads see the latest committed assignment.
|
||||
// A missing singleton fails closed. Deadlocks roll back the losing transaction.
|
||||
func CheckPurchaseOrderNumber(tx *gorm.DB, taskID uint64, orderNo string) error {
|
||||
var setting PurchaseRuleSetting
|
||||
if err := tx.Session(&gorm.Session{NewDB: true}).Clauses(clause.Locking{Strength: "UPDATE"}).First(&setting, 1).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
var others []PurchaseTask
|
||||
if err := tx.Session(&gorm.Session{NewDB: true}).Select("id").Clauses(clause.Locking{Strength: "UPDATE"}).Where("pdd_order_no = ? AND id <> ?", orderNo, taskID).Find(&others).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if len(others) != 0 {
|
||||
return ErrPurchaseOrderNumberUsed
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,202 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/app/goauto/purchasecontract"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
"gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
const (
|
||||
MaxOrderBackfillItems = 50
|
||||
CodeBackfillSuffix = "PURCHASE_BACKFILL_SUFFIX_INVALID"
|
||||
CodeBackfillDevice = "PURCHASE_BACKFILL_DEVICE_MISMATCH"
|
||||
CodeBackfillOrderConflict = "PURCHASE_BACKFILL_ORDER_CONFLICT"
|
||||
CodeBackfillBatchConflict = "PURCHASE_BACKFILL_BATCH_CONFLICT"
|
||||
CodeBackfillOrderUsed = "PURCHASE_BACKFILL_ORDER_ALREADY_USED"
|
||||
)
|
||||
|
||||
type OrderBackfillRequest struct {
|
||||
RequestID string `json:"requestId"`
|
||||
Items []OrderBackfillItem `json:"items"`
|
||||
}
|
||||
|
||||
type OrderBackfillItem struct {
|
||||
AddressSuffix string `json:"addressSuffix"`
|
||||
PDDOrderNo string `json:"pddOrderNo"`
|
||||
// A string keeps an invalid page timestamp local to this item.
|
||||
OrderSubmittedAt *string `json:"orderSubmittedAt,omitempty"`
|
||||
}
|
||||
|
||||
type OrderBackfillResult struct {
|
||||
Index int `json:"index"`
|
||||
TaskID uint64 `json:"taskId,omitempty"`
|
||||
Result string `json:"result"`
|
||||
Code string `json:"code"`
|
||||
Status string `json:"status,omitempty"`
|
||||
StatusVersion uint64 `json:"statusVersion,omitempty"`
|
||||
PDDOrderNo *string `json:"pddOrderNo,omitempty"`
|
||||
OrderSubmittedAt *time.Time `json:"orderSubmittedAt,omitempty"`
|
||||
TimeSource string `json:"timeSource,omitempty"`
|
||||
Retryable bool `json:"retryable"`
|
||||
}
|
||||
|
||||
type OrderBackfillResponse struct {
|
||||
RequestID string `json:"requestId"`
|
||||
Items []OrderBackfillResult `json:"items"`
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
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)
|
||||
for i, item := range req.Items {
|
||||
id, err := purchasecontract.ParseAddressSuffix(item.AddressSuffix)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
ids[i] = id
|
||||
if previous, ok := orders[id]; ok && previous != item.PDDOrderNo {
|
||||
conflicts[id] = true
|
||||
}
|
||||
orders[id] = item.PDDOrderNo
|
||||
}
|
||||
out.Items = make([]OrderBackfillResult, len(req.Items))
|
||||
for i, item := range req.Items {
|
||||
r := OrderBackfillResult{Index: i, TaskID: ids[i], Result: "failed"}
|
||||
if ids[i] == 0 {
|
||||
r.Code = CodeBackfillSuffix
|
||||
} else {
|
||||
r = s.backfillOrder(ctx, d.ID, ids[i], req.RequestID, item, conflicts[ids[i]])
|
||||
r.Index = i
|
||||
}
|
||||
out.Items[i] = r
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (s *Service) backfillOrder(ctx context.Context, deviceID, 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.
|
||||
db := s.DB.Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent)}).WithContext(ctx)
|
||||
err := db.Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&task, taskID).Error; err != nil {
|
||||
return purchaseNotFound(err)
|
||||
}
|
||||
if task.DeviceID == nil || *task.DeviceID != deviceID {
|
||||
return fail(CodeBackfillDevice, "任务不属于当前设备")
|
||||
}
|
||||
if batchConflict {
|
||||
return fail(CodeBackfillBatchConflict, "同批任务有不同订单号")
|
||||
}
|
||||
if task.ExecutionMode != models.PurchaseExecutionModeLive || (task.Status != models.PurchaseTaskStatusOrderResultUnknown && task.Status != models.PurchaseTaskStatusOrderCreated) {
|
||||
return fail(CodeStateConflict, "当前任务不允许回填")
|
||||
}
|
||||
if item.PDDOrderNo == "" || strings.TrimSpace(item.PDDOrderNo) != item.PDDOrderNo || utf8.RuneCountInString(item.PDDOrderNo) > 100 || strings.ContainsAny(item.PDDOrderNo, "\r\n\t") {
|
||||
return fail(CodeInvalidRequest, "订单号无效")
|
||||
}
|
||||
if task.PDDOrderNo != nil && *task.PDDOrderNo != "" && *task.PDDOrderNo != item.PDDOrderNo {
|
||||
return fail(CodeBackfillOrderConflict, "已有不同订单号")
|
||||
}
|
||||
// The shared model guard also protects manual resolution and late results.
|
||||
if err := models.CheckPurchaseOrderNumber(tx, taskID, item.PDDOrderNo); err != nil {
|
||||
return err
|
||||
}
|
||||
if task.Status == models.PurchaseTaskStatusOrderCreated {
|
||||
if task.PDDOrderNo == nil || *task.PDDOrderNo != item.PDDOrderNo {
|
||||
return fail(CodeStateConflict, "已创建订单缺少匹配订单号")
|
||||
}
|
||||
r.Result, r.Code = "already_backfilled", "ALREADY_BACKFILLED"
|
||||
return nil
|
||||
}
|
||||
var submitted time.Time
|
||||
source := "page"
|
||||
if item.OrderSubmittedAt != nil {
|
||||
var err error
|
||||
submitted, err = time.Parse(time.RFC3339Nano, *item.OrderSubmittedAt)
|
||||
if err != nil || submitted.IsZero() || submitted.Year() < 1000 || submitted.Year() > 9999 {
|
||||
return fail(CodeOrderTimeInvalid, "下单时间必须为 RFC3339")
|
||||
}
|
||||
} else {
|
||||
if task.IrreversibleAt == nil || task.IrreversibleAt.IsZero() {
|
||||
return fail(CodeOrderTimeMissing, "下单时间和不可逆时间均缺失")
|
||||
}
|
||||
submitted, source = *task.IrreversibleAt, "irreversible_at"
|
||||
}
|
||||
submitted = submitted.UTC()
|
||||
task.PDDOrderNo, task.OrderSubmittedAt = &item.PDDOrderNo, &submitted
|
||||
if err := task.SetStatus(models.PurchaseTaskStatusOrderCreated); err != nil {
|
||||
return internal(err)
|
||||
}
|
||||
task.StatusVersion++
|
||||
task.StatusChangedAt = s.Now()
|
||||
task.ErrorCode, task.ErrorMessage = nil, nil
|
||||
task.LeaseExpiresAt = nil
|
||||
// Reuse the existing resolution request slot. Scope a batch UUID to a
|
||||
// task, and retain provenance without a schema change or replay cache.
|
||||
marker := "backfill:" + source + ":" + uuid.NewSHA1(uuid.NameSpaceOID, []byte(requestID+":"+item.AddressSuffix)).String()
|
||||
task.UnknownResolveRequestID = &marker
|
||||
if err := tx.Save(&task).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
r.Result, r.Code = "backfilled", "BACKFILLED"
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
r.Result, r.Code = "failed", CodeInternal
|
||||
r.Retryable = true
|
||||
var se *ServiceError
|
||||
if errors.As(err, &se) {
|
||||
r.Code, r.Retryable = se.Code, se.Retryable
|
||||
}
|
||||
if errors.Is(err, models.ErrPurchaseOrderNumberUsed) {
|
||||
r.Code, r.Retryable = CodeBackfillOrderUsed, false
|
||||
}
|
||||
if r.Code == CodeBackfillBatchConflict || r.Code == CodeBackfillOrderConflict || r.Code == CodeBackfillOrderUsed {
|
||||
r.Result = "conflict"
|
||||
}
|
||||
}
|
||||
// Return only this device's committed facts, including on a rejected item.
|
||||
// Never return in-memory changes from a rolled back transaction.
|
||||
saved := task
|
||||
readable := err == nil
|
||||
if !readable {
|
||||
saved = models.PurchaseTask{}
|
||||
readable = db.Where("id = ? AND device_id = ?", taskID, deviceID).First(&saved).Error == nil
|
||||
}
|
||||
if readable {
|
||||
r.Status, r.StatusVersion = saved.Status, saved.StatusVersion
|
||||
r.PDDOrderNo, r.OrderSubmittedAt = saved.PDDOrderNo, saved.OrderSubmittedAt
|
||||
if saved.OrderSubmittedAt != nil {
|
||||
r.TimeSource = "existing_unknown"
|
||||
if saved.UnknownResolveRequestID != nil {
|
||||
if strings.HasPrefix(*saved.UnknownResolveRequestID, "backfill:page:") {
|
||||
r.TimeSource = "page"
|
||||
}
|
||||
if strings.HasPrefix(*saved.UnknownResolveRequestID, "backfill:irreversible_at:") {
|
||||
r.TimeSource = "irreversible_at"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return r
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
func (h Handler) BackfillOrders(c *gin.Context) {
|
||||
var req OrderBackfillRequest
|
||||
if !decode(c, &req) {
|
||||
return
|
||||
}
|
||||
s, ok := h.service(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
out, err := s.BackfillOrders(c.Request.Context(), 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": out})
|
||||
}
|
||||
@@ -0,0 +1,474 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/app/goauto/purchasecontract"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func backfillTask(t *testing.T, db *gorm.DB, f fixture, status string) models.PurchaseTask {
|
||||
t.Helper()
|
||||
now := testService(db).Now()
|
||||
task := models.PurchaseTask{TaskType: models.PurchaseTaskTypeStock, ExecutionMode: models.PurchaseExecutionModeLive,
|
||||
Status: status, DeviceID: &f.device.ID, PDDProductID: f.pdd.ID, Quantity: 1, Currency: "CNY",
|
||||
CreateRequestID: uuid.NewString(), RuleSnapshot: string(purchasecontract.DefaultLiveRule()),
|
||||
SpecDecisionSnapshot: `{}`, RequiredCapabilitiesJSON: `[]`, IrreversibleAt: &now,
|
||||
ErrorCode: strptr("ORIGINAL_ERROR"), ErrorMessage: strptr("original failure")}
|
||||
if err := db.Create(&task).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return task
|
||||
}
|
||||
|
||||
func strptr(s string) *string { return &s }
|
||||
|
||||
func backfillItem(id uint64, order string) OrderBackfillItem {
|
||||
return OrderBackfillItem{AddressSuffix: purchasecontract.AddressSuffix(id), PDDOrderNo: order}
|
||||
}
|
||||
|
||||
func runBackfill(t *testing.T, s *Service, token, requestID string, items ...OrderBackfillItem) []OrderBackfillResult {
|
||||
t.Helper()
|
||||
out, err := s.BackfillOrders(context.Background(), OrderBackfillRequest{RequestID: requestID, Items: items}, token)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(out.Items) != len(items) || out.RequestID != requestID {
|
||||
t.Fatalf("bad envelope: %+v", out)
|
||||
}
|
||||
return out.Items
|
||||
}
|
||||
|
||||
func loadBackfillTask(t *testing.T, db *gorm.DB, id uint64) models.PurchaseTask {
|
||||
t.Helper()
|
||||
var task models.PurchaseTask
|
||||
if err := db.First(&task, id).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return task
|
||||
}
|
||||
|
||||
func TestOrderBackfillMixedBatchAndReplay(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
a := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
a.TaskType, a.SYBProductID = models.PurchaseTaskTypeSYBOrder, &f.syb.ID
|
||||
if err := db.Save(&a).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
b := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
c := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
if err := db.Model(&c).Update("irreversible_at", nil).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
page := backfillItem(b.ID, "ORDER-B")
|
||||
page.OrderSubmittedAt = strptr("2026-09-08T20:30:00+08:00")
|
||||
rid := uuid.NewString()
|
||||
items := []OrderBackfillItem{backfillItem(a.ID, "ORDER-A"), {AddressSuffix: "_cg0", PDDOrderNo: "bad"}, page, backfillItem(c.ID, "ORDER-C"), backfillItem(99999, "missing")}
|
||||
results := runBackfill(t, s, f.token, rid, items...)
|
||||
want := []string{"BACKFILLED", CodeBackfillSuffix, "BACKFILLED", CodeOrderTimeMissing, CodeTaskNotFound}
|
||||
for i, r := range results {
|
||||
if r.Code != want[i] || r.Index != i {
|
||||
t.Fatalf("item %d: %+v", i, r)
|
||||
}
|
||||
}
|
||||
if results[0].TimeSource != "irreversible_at" || !results[0].OrderSubmittedAt.Equal(*a.IrreversibleAt) {
|
||||
t.Fatalf("fallback: %+v", results[0])
|
||||
}
|
||||
if results[2].TimeSource != "page" || results[2].OrderSubmittedAt.Format(time.RFC3339) != "2026-09-08T12:30:00Z" {
|
||||
t.Fatalf("page: %+v", results[2])
|
||||
}
|
||||
saved := loadBackfillTask(t, db, a.ID)
|
||||
if saved.StatusVersion != a.StatusVersion+1 || saved.ErrorCode != nil || saved.ErrorMessage != nil || saved.DeviceRunSlot != nil || saved.AccountRunSlot != nil || saved.ActiveSlot == nil || saved.Status != models.PurchaseTaskStatusOrderCreated {
|
||||
t.Fatalf("state metadata: %+v", saved)
|
||||
}
|
||||
if saved.PaymentReviewStatus != a.PaymentReviewStatus || saved.LogisticsStatus != a.LogisticsStatus || saved.WritebackStatus != a.WritebackStatus || saved.RuleSnapshot != a.RuleSnapshot {
|
||||
t.Fatal("unrelated business facts changed")
|
||||
}
|
||||
for _, replayID := range []string{rid, uuid.NewString()} {
|
||||
item := items[0]
|
||||
item.OrderSubmittedAt = strptr("2026-09-09T00:00:00Z")
|
||||
r := runBackfill(t, s, f.token, replayID, item)[0]
|
||||
if r.Result != "already_backfilled" || r.TimeSource != "irreversible_at" {
|
||||
t.Fatalf("replay: %+v", r)
|
||||
}
|
||||
if got := loadBackfillTask(t, db, a.ID); !reflect.DeepEqual(saved, got) {
|
||||
t.Fatal("replay changed persisted task")
|
||||
}
|
||||
}
|
||||
if got := loadBackfillTask(t, db, c.ID); got.PDDOrderNo != nil || got.StatusVersion != c.StatusVersion {
|
||||
t.Fatal("missing time wrote data")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillRejectsOwnershipStatesAndInvalidTime(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
for _, status := range []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted, models.PurchaseTaskStatusSpecProbePending, models.PurchaseTaskStatusFailed, models.PurchaseTaskStatusCancelled, models.PurchaseTaskStatusRehearsalCompleted} {
|
||||
task := backfillTask(t, db, f, status)
|
||||
before := loadBackfillTask(t, db, task.ID)
|
||||
r := runBackfill(t, s, f.token, uuid.NewString(), backfillItem(task.ID, "ORDER"))[0]
|
||||
if r.Code != CodeStateConflict {
|
||||
t.Fatalf("%s: %+v", status, r)
|
||||
}
|
||||
if got := loadBackfillTask(t, db, task.ID); !reflect.DeepEqual(got, before) {
|
||||
t.Fatal("rejection wrote data")
|
||||
}
|
||||
// Release the fixture's device slot before testing the next running state.
|
||||
if err := task.SetStatus(models.PurchaseTaskStatusCancelled); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Save(&task).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
task := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
if err := db.Model(&task).Update("device_id", nil).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
r := runBackfill(t, s, f.token, uuid.NewString(), backfillItem(task.ID, "ORDER"))[0]
|
||||
if r.Code != CodeBackfillDevice || r.Status != "" || r.PDDOrderNo != nil {
|
||||
t.Fatalf("ownership leaked: %+v", r)
|
||||
}
|
||||
other, err := device.NewService(db).Register(context.Background(), device.RegisterRequest{RequestID: uuid.NewString(), InstallID: uuid.NewString(), Name: "Other", Manufacturer: "Test", Model: "Test", AndroidVersion: "15", AgentVersion: "1", PDDVersion: "7", Capabilities: liveCaps()}, "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Model(&task).Update("device_id", other.DeviceID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if r := runBackfill(t, s, f.token, uuid.NewString(), backfillItem(task.ID, "ORDER"))[0]; r.Code != CodeBackfillDevice {
|
||||
t.Fatalf("cross device: %+v", r)
|
||||
}
|
||||
if err := db.Model(&task).Updates(map[string]any{"device_id": f.device.ID, "execution_mode": models.PurchaseExecutionModeRehearsal}).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if r := runBackfill(t, s, f.token, uuid.NewString(), backfillItem(task.ID, "ORDER"))[0]; r.Code != CodeStateConflict {
|
||||
t.Fatalf("rehearsal: %+v", r)
|
||||
}
|
||||
if err := db.Model(&task).Update("execution_mode", models.PurchaseExecutionModeLive).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, raw := range []string{"", "2026-09-08 12:00:00", "0001-01-01T00:00:00Z", "garbage"} {
|
||||
item := backfillItem(task.ID, "ORDER")
|
||||
item.OrderSubmittedAt = &raw
|
||||
if r := runBackfill(t, s, f.token, uuid.NewString(), item)[0]; r.Code != CodeOrderTimeInvalid {
|
||||
t.Fatalf("invalid time: %+v", r)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillConflictsNeverOverwrite(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
a := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
b := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
rid := uuid.NewString()
|
||||
r := runBackfill(t, s, f.token, rid, backfillItem(a.ID, "A"), backfillItem(a.ID, "B"), backfillItem(b.ID, "B"))
|
||||
if r[0].Code != CodeBackfillBatchConflict || r[1].Code != CodeBackfillBatchConflict || r[2].Code != "BACKFILLED" {
|
||||
t.Fatalf("batch: %+v", r)
|
||||
}
|
||||
r = runBackfill(t, s, f.token, rid, backfillItem(a.ID, "B"), backfillItem(b.ID, "C"))
|
||||
if r[0].Code != CodeBackfillOrderUsed || r[1].Code != CodeBackfillOrderConflict {
|
||||
t.Fatalf("changed requestId payload bypassed checks: %+v", r)
|
||||
}
|
||||
if got := loadBackfillTask(t, db, b.ID); *got.PDDOrderNo != "B" || got.StatusVersion != b.StatusVersion+1 {
|
||||
t.Fatal("conflict overwrote")
|
||||
}
|
||||
if got := loadBackfillTask(t, db, a.ID); got.PDDOrderNo != nil {
|
||||
t.Fatal("conflict wrote data")
|
||||
}
|
||||
// Even an unknown task with an existing conflicting value must preserve it.
|
||||
a.PDDOrderNo = strptr("OLD")
|
||||
if err := db.Save(&a).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if r := runBackfill(t, s, f.token, uuid.NewString(), backfillItem(a.ID, "NEW"))[0]; r.Code != CodeBackfillOrderConflict {
|
||||
t.Fatalf("unknown existing: %+v", r)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillConcurrentResolveUnknown(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
// SQLite serializes transactions through one connection. These concurrent
|
||||
// service calls verify both winner orders; they do not certify MySQL locks.
|
||||
sqlDB, _ := db.DB()
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
for i := 0; i < 12; i++ {
|
||||
task := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
start := make(chan struct{})
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(2)
|
||||
var out OrderBackfillResponse
|
||||
var backErr, manualErr error
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
out, backErr = s.BackfillOrders(context.Background(), OrderBackfillRequest{RequestID: uuid.NewString(), Items: []OrderBackfillItem{backfillItem(task.ID, "BACK-"+purchasecontract.AddressSuffix(task.ID))}}, f.token)
|
||||
}()
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
_, _, manualErr = s.ResolveUnknown(context.Background(), task.ID, ManualRequest{RequestID: uuid.NewString(), OperatorID: 1, Status: models.PurchaseTaskStatusOrderCreated, PDDOrderNo: "MANUAL-" + purchasecontract.AddressSuffix(task.ID), OrderSubmittedAt: task.IrreversibleAt})
|
||||
}()
|
||||
close(start)
|
||||
wg.Wait()
|
||||
if backErr != nil {
|
||||
t.Fatal(backErr)
|
||||
}
|
||||
got := loadBackfillTask(t, db, task.ID)
|
||||
if got.StatusVersion != task.StatusVersion+1 || got.Status != models.PurchaseTaskStatusOrderCreated {
|
||||
t.Fatal("competing writes changed version twice")
|
||||
}
|
||||
if manualErr == nil {
|
||||
if out.Items[0].Code != CodeBackfillOrderConflict || !strings.HasPrefix(*got.PDDOrderNo, "MANUAL-") {
|
||||
t.Fatalf("manual winner: %+v", out)
|
||||
}
|
||||
} else if code(manualErr) != CodeStateConflict || out.Items[0].Code != "BACKFILLED" || !strings.HasPrefix(*got.PDDOrderNo, "BACK-") {
|
||||
t.Fatalf("backfill winner: %+v %v", out, manualErr)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillConcurrentLateResultAndOtherTask(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
sqlDB, _ := db.DB()
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
a := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
attempt := models.PurchaseTaskAttempt{TaskID: a.ID, AttemptID: uuid.NewString(), AttemptNumber: 1, Phase: models.PurchaseAttemptPhasePurchase, Status: models.PurchaseAttemptStatusFailed, DeviceID: &f.device.ID, RuleSnapshotHash: purchaseRuleSnapshotHash(a.RuleSnapshot), SpecDecisionSnapshot: `{}`}
|
||||
if err := db.Omit("Task").Create(&attempt).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.First(&attempt, attempt.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
start := make(chan struct{})
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(2)
|
||||
var out OrderBackfillResponse
|
||||
var backErr, lateErr error
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
out, backErr = s.BackfillOrders(context.Background(), OrderBackfillRequest{RequestID: uuid.NewString(), Items: []OrderBackfillItem{backfillItem(a.ID, "BACK")}}, f.token)
|
||||
}()
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
_, lateErr = s.SubmitResult(context.Background(), a.ID, ResultRequest{RequestID: uuid.NewString(), TaskAttemptID: attempt.AttemptID, ResultType: "order_created", PDDOrderNo: "LATE", OrderSubmittedAt: a.IrreversibleAt}, f.token)
|
||||
}()
|
||||
close(start)
|
||||
wg.Wait()
|
||||
if backErr != nil || out.Items[0].Code != "BACKFILLED" || code(lateErr) != CodeStateConflict {
|
||||
t.Fatalf("late race: %+v %v %v", out, backErr, lateErr)
|
||||
}
|
||||
var savedAttempt models.PurchaseTaskAttempt
|
||||
if err := db.First(&savedAttempt, attempt.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !reflect.DeepEqual(savedAttempt, attempt) {
|
||||
t.Fatal("backfill rewrote attempt")
|
||||
}
|
||||
b := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
_, _, err := s.ResolveUnknown(context.Background(), b.ID, ManualRequest{RequestID: uuid.NewString(), OperatorID: 1, Status: models.PurchaseTaskStatusOrderCreated, PDDOrderNo: "BACK", OrderSubmittedAt: b.IrreversibleAt})
|
||||
if err == nil {
|
||||
t.Fatal("manual path assigned another task's order")
|
||||
}
|
||||
if got := loadBackfillTask(t, db, b.ID); got.PDDOrderNo != nil || got.StatusVersion != b.StatusVersion {
|
||||
t.Fatal("other task changed on conflict")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillHTTPBoundary(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
task := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
gin.SetMode(gin.TestMode)
|
||||
r := gin.New()
|
||||
r.POST("/order-backfill", device.RequireAgentHTTPS(false, false), (Handler{DB: db}).BackfillOrders)
|
||||
body, _ := json.Marshal(OrderBackfillRequest{RequestID: uuid.NewString(), Items: []OrderBackfillItem{backfillItem(task.ID, "HTTP")}})
|
||||
for _, test := range []struct {
|
||||
body, token string
|
||||
status int
|
||||
}{
|
||||
{string(body), "", http.StatusUnauthorized},
|
||||
{`{"requestId":"bad","items":[]}`, f.token, http.StatusUnprocessableEntity},
|
||||
{`{"requestId":"x","address":"forbidden"}`, f.token, http.StatusUnprocessableEntity},
|
||||
{string(body), f.token, http.StatusOK},
|
||||
} {
|
||||
req := httptest.NewRequest(http.MethodPost, "/order-backfill", strings.NewReader(test.body))
|
||||
req.Header.Set("Authorization", "Bearer "+test.token)
|
||||
w := httptest.NewRecorder()
|
||||
r.ServeHTTP(w, req)
|
||||
if w.Code != test.status {
|
||||
t.Fatalf("HTTP %d: %s", w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
_, err := testService(db).BackfillOrders(context.Background(), OrderBackfillRequest{RequestID: uuid.NewString(), Items: make([]OrderBackfillItem, 51)}, f.token)
|
||||
if code(err) != CodeInvalidRequest {
|
||||
t.Fatalf("batch limit: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillConcurrentSameOrderDifferentTasks(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
sqlDB, _ := db.DB()
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
a := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
b := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
start := make(chan struct{})
|
||||
results := make(chan OrderBackfillResponse, 2)
|
||||
errors := make(chan error, 2)
|
||||
for _, id := range []uint64{a.ID, b.ID} {
|
||||
go func(id uint64) {
|
||||
<-start
|
||||
out, err := s.BackfillOrders(context.Background(), OrderBackfillRequest{RequestID: uuid.NewString(), Items: []OrderBackfillItem{backfillItem(id, "SAME")}}, f.token)
|
||||
results <- out
|
||||
errors <- err
|
||||
}(id)
|
||||
}
|
||||
close(start)
|
||||
codes := make(map[string]int)
|
||||
for i := 0; i < 2; i++ {
|
||||
out := <-results
|
||||
if err := <-errors; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
codes[out.Items[0].Code]++
|
||||
}
|
||||
if codes["BACKFILLED"] != 1 || codes[CodeBackfillOrderUsed] != 1 {
|
||||
t.Fatalf("concurrent assignments: %+v", codes)
|
||||
}
|
||||
var count int64
|
||||
if err := db.Model(&models.PurchaseTask{}).Where("pdd_order_no = ?", "SAME").Count(&count).Error; err != nil || count != 1 {
|
||||
t.Fatalf("duplicate order: %d %v", count, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillRejectsLateAssignmentOfSameOrder(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
a := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
if r := runBackfill(t, s, f.token, uuid.NewString(), backfillItem(a.ID, "SHARED"))[0]; r.Code != "BACKFILLED" {
|
||||
t.Fatal(r)
|
||||
}
|
||||
b := backfillTask(t, db, f, models.PurchaseTaskStatusOrderSubmitStarted)
|
||||
lease := s.Now().Add(time.Minute)
|
||||
b.LeaseExpiresAt = &lease
|
||||
if err := db.Save(&b).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
attempt := models.PurchaseTaskAttempt{TaskID: b.ID, AttemptID: uuid.NewString(), AttemptNumber: 1, Phase: models.PurchaseAttemptPhasePurchase, Status: models.PurchaseAttemptStatusRunning, DeviceID: &f.device.ID, RuleSnapshotHash: purchaseRuleSnapshotHash(b.RuleSnapshot), SpecDecisionSnapshot: `{}`}
|
||||
if err := db.Omit("Task").Create(&attempt).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err := s.SubmitResult(context.Background(), b.ID, ResultRequest{RequestID: uuid.NewString(), TaskAttemptID: attempt.AttemptID, ResultType: "order_created", PDDOrderNo: "SHARED", OrderSubmittedAt: b.IrreversibleAt}, f.token)
|
||||
if err == nil {
|
||||
t.Fatal("old result path assigned duplicate order")
|
||||
}
|
||||
got := loadBackfillTask(t, db, b.ID)
|
||||
if got.StatusVersion != b.StatusVersion || got.PDDOrderNo != nil {
|
||||
t.Fatal("old result transaction was not rolled back")
|
||||
}
|
||||
var saved models.PurchaseTaskAttempt
|
||||
if err := db.First(&saved, attempt.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if saved.Status != models.PurchaseAttemptStatusRunning || saved.ResultRequestID != nil {
|
||||
t.Fatal("attempt result survived rolled back assignment")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillHTTPTransportPolicy(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
for _, allow := range []string{"false", "true"} {
|
||||
t.Setenv("GOAUTO_ALLOW_INSECURE_AGENT_HTTP", allow)
|
||||
r := gin.New()
|
||||
r.POST("/order-backfill", device.RequireAgentHTTPS(true, false), (Handler{}).BackfillOrders)
|
||||
w := httptest.NewRecorder()
|
||||
r.ServeHTTP(w, httptest.NewRequest(http.MethodPost, "/order-backfill", strings.NewReader(`{}`)))
|
||||
if allow == "false" && w.Code != http.StatusUpgradeRequired {
|
||||
t.Fatalf("HTTPS bypass: %d", w.Code)
|
||||
}
|
||||
if allow == "true" && w.Code == http.StatusUpgradeRequired {
|
||||
t.Fatal("HTTP compatibility broken")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderBackfillMultiConnectionResolveRace(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB.SetMaxOpenConns(4)
|
||||
task := backfillTask(t, db, f, models.PurchaseTaskStatusOrderResultUnknown)
|
||||
start := make(chan struct{})
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(2)
|
||||
var back OrderBackfillResponse
|
||||
var backErr, manualErr error
|
||||
req := OrderBackfillRequest{RequestID: uuid.NewString(), Items: []OrderBackfillItem{backfillItem(task.ID, "BACK")}}
|
||||
go func() { defer wg.Done(); <-start; back, backErr = s.BackfillOrders(context.Background(), req, f.token) }()
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
_, _, manualErr = s.ResolveUnknown(context.Background(), task.ID, ManualRequest{RequestID: uuid.NewString(), OperatorID: 1, Status: models.PurchaseTaskStatusOrderCreated, PDDOrderNo: "MANUAL", OrderSubmittedAt: task.IrreversibleAt})
|
||||
}()
|
||||
close(start)
|
||||
wg.Wait()
|
||||
// SQLite returns table-lock errors rather than waiting on FOR UPDATE.
|
||||
// Only that documented DB contention or a domain conflict is acceptable;
|
||||
// after the competing calls finish, replay must converge without overwrite.
|
||||
if backErr != nil && !strings.Contains(backErr.Error(), "locked") {
|
||||
t.Fatal(backErr)
|
||||
}
|
||||
if manualErr != nil && code(manualErr) != CodeStateConflict && !strings.Contains(manualErr.Error(), "locked") {
|
||||
t.Fatal(manualErr)
|
||||
}
|
||||
if backErr == nil && back.Items[0].Code != "BACKFILLED" && back.Items[0].Code != CodeBackfillOrderConflict && !(back.Items[0].Code == CodeInternal && back.Items[0].Retryable) {
|
||||
t.Fatalf("unexpected race result: %+v", back)
|
||||
}
|
||||
before := loadBackfillTask(t, db, task.ID)
|
||||
replay := runBackfill(t, s, f.token, req.RequestID, req.Items...)[0]
|
||||
after := loadBackfillTask(t, db, task.ID)
|
||||
if before.PDDOrderNo != nil && !reflect.DeepEqual(before, after) {
|
||||
t.Fatal("replay overwrote the concurrent winner")
|
||||
}
|
||||
if after.StatusVersion != task.StatusVersion+1 || after.Status != models.PurchaseTaskStatusOrderCreated {
|
||||
t.Fatal("race did not converge to a single transition")
|
||||
}
|
||||
if manualErr == nil {
|
||||
if *after.PDDOrderNo != "MANUAL" || replay.Code != CodeBackfillOrderConflict {
|
||||
t.Fatal("manual winner overwritten")
|
||||
}
|
||||
} else if *after.PDDOrderNo != "BACK" || (replay.Code != "BACKFILLED" && replay.Code != "ALREADY_BACKFILLED") {
|
||||
t.Fatalf("backfill did not converge: %+v", replay)
|
||||
}
|
||||
}
|
||||
@@ -18,6 +18,7 @@ func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) {
|
||||
agent := engine.Group("/api/agent/v1/purchase-tasks").Use(device.RequireAgentHTTPS(config.ApplicationConfig.Mode == "prod", trust))
|
||||
agent.GET("", h.AgentHistory)
|
||||
agent.GET("/next", h.Next)
|
||||
agent.POST("/order-backfill", h.BackfillOrders)
|
||||
agent.GET("/:taskId", h.AgentHistoryDetail)
|
||||
agent.POST("/:taskId/retry", h.AgentRetry)
|
||||
agent.POST("/:taskId/reset", h.AgentReset)
|
||||
|
||||
@@ -0,0 +1,17 @@
|
||||
package purchasecontract
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestParseAddressSuffix(t *testing.T) {
|
||||
for _, id := range []uint64{7, 72, ^uint64(0)} {
|
||||
got, err := ParseAddressSuffix(AddressSuffix(id))
|
||||
if err != nil || got != id {
|
||||
t.Fatalf("id=%d got=%d err=%v", id, got, err)
|
||||
}
|
||||
}
|
||||
for _, raw := range []string{"", "_cg", "_cg0", "_cg00", "_cg07", "_cg+7", "_cg-7", "_cg18446744073709551616", "_CG7", "_cg7x", "_cg7_cg72", "address_cg7", " _cg7", "_cg7 ", "_cg7", "_cg7\n"} {
|
||||
if id, err := ParseAddressSuffix(raw); err == nil || id != 0 {
|
||||
t.Errorf("accepted %q: %d", raw, id)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
"math"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"unicode/utf8"
|
||||
)
|
||||
@@ -360,6 +361,18 @@ func RequiredCapabilities(rule RuleSnapshot) []string {
|
||||
|
||||
func AddressSuffix(taskID uint64) string { return fmt.Sprintf("_cg%d", taskID) }
|
||||
|
||||
// ParseAddressSuffix accepts only the exact canonical suffix, never an address.
|
||||
func ParseAddressSuffix(suffix string) (uint64, error) {
|
||||
if !strings.HasPrefix(suffix, "_cg") {
|
||||
return 0, errors.New("invalid address suffix")
|
||||
}
|
||||
id, err := strconv.ParseUint(strings.TrimPrefix(suffix, "_cg"), 10, 64)
|
||||
if err != nil || id == 0 || AddressSuffix(id) != suffix {
|
||||
return 0, errors.New("invalid address suffix")
|
||||
}
|
||||
return id, nil
|
||||
}
|
||||
|
||||
func ensureEOF(decoder *json.Decoder) error {
|
||||
var extra any
|
||||
if err := decoder.Decode(&extra); err != io.EOF {
|
||||
|
||||
Reference in New Issue
Block a user