feat(#34): implement purchase task state machine
This commit is contained in:
@@ -2,7 +2,7 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Project-Profile
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Project-Profile.-
|
||||
wiki_revision: f9245c856a45c10dfbae7fe32b365ca11b008782
|
||||
wiki_revision: 5da7f2250f851caab6d0d1713f869b619ed55eb2
|
||||
synchronized_at: 2026-08-20T08:47:11Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
@@ -70,4 +70,4 @@ synchronized_at: 2026-08-20T08:47:11Z
|
||||
|
||||
#40、#41、#49~#52 已于 2026-08-20 通过用户验收,虾皮/SYB 商品档案、店铺过滤、异步同步记录、列表布局和 MySQL 同步缺陷修复均已完成。
|
||||
|
||||
#33 已建立采购任务、任务尝试和可选 PDD 账号引用的数据契约,完成规则能力隔离与 MySQL 8.4 迁移验证,等待用户验收。HTTP 状态机、Admin 页面和 Android 执行仍属于后续工单;当前实现不会创建 PDD 订单。
|
||||
#33 已于 2026-08-20 通过用户验收。#34 已完成服务端采购任务创建、设备/可选账号租约、两趟规格探测、attempt 结果幂等、订单结果未知和人工处置状态机,当前等待用户验收;Admin 页面和 Android 执行动作仍属于后续工单,本次没有创建 PDD 订单。
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Architecture-and-Code-Map
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.-
|
||||
wiki_revision: f9245c856a45c10dfbae7fe32b365ca11b008782
|
||||
wiki_revision: 5da7f2250f851caab6d0d1713f869b619ed55eb2
|
||||
synchronized_at: 2026-08-20T08:47:11Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
@@ -118,6 +118,7 @@ Android Portal/Agent
|
||||
| SYB 店铺准入、发现与过滤 | `server/app/goauto/sybshop/`、`server/app/goauto/sybimport/`;迁移 `server/cmd/migrate/migration/version-local/1786700900000_syb_shop.go` |
|
||||
| SYB 后台导入任务、进度、单任务互斥与启动恢复 | `server/app/goauto/sybimport/sync_run.go`、`sync_run_handler.go`;表 `syb_sync_run`,迁移 `server/cmd/migrate/migration/version-local/1786701000000_syb_sync_run.go` |
|
||||
| 采购任务数据与类型化规则契约 | `server/app/goauto/models/purchase.go`、`server/app/goauto/purchasecontract/`;迁移 `server/cmd/migrate/migration/version-local/1786701100000_purchase_contract.go` |
|
||||
| 采购任务创建、租约、attempt 幂等和人工处置状态机 | `server/app/goauto/purchase/`;追加迁移 `server/cmd/migrate/migration/version-local/1786701200000_purchase_state_machine.go` |
|
||||
| 任务领取、结果、重置与删除 | `server/app/goauto/task/` |
|
||||
| 管理端基线 | `web/`(go-admin-ui v3.0.0) |
|
||||
| 管理端闭环页面 | `web/src/views/goauto/` |
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Business-Rules-and-Glossary
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.-
|
||||
wiki_revision: f9245c856a45c10dfbae7fe32b365ca11b008782
|
||||
wiki_revision: 5da7f2250f851caab6d0d1713f869b619ed55eb2
|
||||
synchronized_at: 2026-08-20T08:47:11Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
@@ -117,6 +117,11 @@ synchronized_at: 2026-08-20T08:47:11Z
|
||||
- 点击创建订单前必须先保存 `order_submit_started` 和不可逆时间;结果不明时进入 `order_result_unknown`,禁止自动再次点击。
|
||||
- Agent 本地 Room/Outbox 负责断网和重启恢复,服务端以 `task_id + task_attempt_id` 幂等接收并保存最终事实。
|
||||
- 人工支付复核只记录 `paid` / `unpaid`;系统不执行或识别支付。快递单号与回填状态属于采购任务,后续物流工单实现。
|
||||
- 采购任务领取时同时占用设备租约和可选 PDD 账号租约;租约过期后才可释放并重新领取。设备还存在采集任务时不能领取采购任务。
|
||||
- 规格映射不完整、PDD 档案为待采集或没有规格时,规则必须具有 `purchase.spec-probe.v1`;第一趟只探测规格并释放租约,服务端固化同一 attempt 的决策后才派发第二趟。无匹配结果明确失败。
|
||||
- Agent 提交的相同 attempt 最终结果只能写入一次;相同请求重放返回原事实,不同内容拒绝覆盖。`order_result_unknown` 不参与自动派发,只能人工解除。
|
||||
- 已创建订单默认禁止再次采购;管理员或采购员可以做一次性重新采购授权,新任务创建成功时在同一事务消耗授权,旧任务和旧订单保留。
|
||||
- 人工回填候选只允许从已支付订单选择;同一 SYB 明细后来选择的订单覆盖旧候选,但不删除旧订单事实。
|
||||
|
||||
## 自动化边界
|
||||
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Android-Agent-API-Contract
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Android-Agent-API-Contract.-
|
||||
wiki_revision: f9245c856a45c10dfbae7fe32b365ca11b008782
|
||||
wiki_revision: 5da7f2250f851caab6d0d1713f869b619ed55eb2
|
||||
synchronized_at: 2026-08-20T08:47:11Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
@@ -369,9 +369,9 @@ POST /api/agent/v1/tasks/{taskId}/fail
|
||||
| `DEVICE_BUSY` | 设备已有活动任务 | 否 |
|
||||
| `DEVICE_CAPABILITY_MISMATCH` | 设备缺少任务规则要求的版本化能力 | 否 |
|
||||
|
||||
## 采购任务共享契约(#33)
|
||||
## 采购任务共享契约(#33、#34)
|
||||
|
||||
本节固定采购域的数据和接口边界;`purchase_task`、`purchase_task_attempt` 与规则校验已经落地,HTTP 路由、状态机服务和 Admin 页面分别由 #34、#35、#44 实现。任何实现都不得扩展为自动支付。
|
||||
本节固定采购域的数据和接口边界;`purchase_task`、`purchase_task_attempt`、规则校验及服务端 HTTP 状态机已落地,Admin 页面和 Android 执行动作分别由后续工单实现。任何实现都不得扩展为自动支付。
|
||||
|
||||
### 任务与快照
|
||||
|
||||
@@ -419,17 +419,19 @@ POST /api/agent/v1/tasks/{taskId}/fail
|
||||
|
||||
演练规则必须包含 `purchase.rehearsal.v1`,并且不能包含 `updateShippingAddress`、`createOrder`、`readOrderResult`。正式规则必须包含 `purchase.live.v1`;改地址和创建订单还分别要求 `purchase.address-update.v1`、`purchase.order-create.v1`。`probeSpecs` 要求 `purchase.spec-probe.v1`。任意模式下,`pay`、名称包含 `payment` 的动作以及未知动作一律拒绝;服务端不下发任意脚本。
|
||||
|
||||
### 管理端接口(由 #34/#35/#44 实现)
|
||||
### 管理端接口
|
||||
|
||||
| 方法 | 路径 | 幂等键 / 说明 |
|
||||
|---|---|---|
|
||||
| `POST` | `/api/admin/v1/purchase-tasks` | `requestId`;单条创建 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/batch` | 批次 `requestId`;逐条成功或失败,不合并任务 |
|
||||
| `GET` | `/api/admin/v1/purchase-tasks` | 分页列表 |
|
||||
| `GET` | `/api/admin/v1/purchase-tasks/{taskId}` | 任务、快照、attempt、订单与物流事实 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/{taskId}/authorize-repurchase` | 一次性授权;创建新任务后自动消耗 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/{taskId}/payment-review` | `paid` 或 `unpaid`,管理员和采购员可操作 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/{taskId}/select-writeback` | 人工选择回填候选;允许明确覆盖旧选择 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/{taskId}/writeback-candidate` | 人工选择回填候选;后选的已支付订单覆盖同一 SYB 明细的旧选择 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/{taskId}/cancel` | 人工取消;执行中和已开始提交订单时禁止直接取消 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/{taskId}/resolve-unknown` | 人工把结果未知解除为已创建订单或已取消 |
|
||||
| `POST` | `/api/admin/v1/purchase-tasks/{taskId}/spec-decision` | 固化第一趟规格探测决策;同一 attempt 不可改写 |
|
||||
|
||||
本工单只提供状态机写接口;列表、详情、批量创建等 Admin 查询与交互接口由 #35 按已确认原型补齐。
|
||||
|
||||
创建请求的价格保护使用整数分:`referenceUnitPriceCent`、`minUnitPriceCent`、`maxUnitPriceCent` 和 `currency`。执行时以 PDD App 实际单价校验;低于最小值或高于最大值均返回普通人可理解的价格越界错误,不考虑优惠券,不以订单总价替代单价判断。
|
||||
|
||||
@@ -440,9 +442,10 @@ POST /api/agent/v1/tasks/{taskId}/fail
|
||||
| `GET` | `/api/agent/v1/purchase-tasks/next` | 返回与设备能力兼容的指定任务或空闲任务 |
|
||||
| `POST` | `/api/agent/v1/purchase-tasks/{taskId}/claim` | `requestId` 原子领取,并建立设备/可选账号租约 |
|
||||
| `POST` | `/api/agent/v1/purchase-tasks/{taskId}/start` | 创建不可变 `taskAttemptId` |
|
||||
| `POST` | `/api/agent/v1/purchase-tasks/{taskId}/attempts/{taskAttemptId}/result` | `requestId`;幂等提交演练、规格探测、订单或失败结果 |
|
||||
| `POST` | `/api/agent/v1/purchase-tasks/{taskId}/order-submit-started` | 创建订单前先落不可逆标记;演练任务永远拒绝 |
|
||||
| `POST` | `/api/agent/v1/purchase-tasks/{taskId}/result` | 请求体携带 `taskAttemptId` 和 `requestId`;幂等提交演练、规格探测、订单或失败结果 |
|
||||
|
||||
结果提交至少关联 `taskId`、`taskAttemptId`、`deviceId`、规则快照哈希和结构化结果。相同 attempt 的重复提交必须返回同一事实;不同内容不得覆盖。慢路径第一趟提交规格后释放设备租约,任务进入 `spec_probe_pending`;服务端固化同一 attempt 的 AI 决策后,第二趟使用新的 attempt 重新派发。
|
||||
结果提交至少关联 `taskId`、`taskAttemptId`、`deviceId`、规则快照哈希和结构化结果。相同 attempt 的相同结果重复提交返回同一事实;不同内容拒绝覆盖。慢路径第一趟提交规格后释放设备与已知账号租约,任务进入 `spec_probe_pending`;服务端固化同一 attempt 的 AI/人工决策后,第二趟使用新的 attempt 重新派发。无匹配结果则明确失败。`order_result_unknown` 只允许管理员或采购员人工解除,永不自动重派。
|
||||
|
||||
Agent Room 只保存恢复执行所需的任务、attempt 和 Outbox;服务端数据库是最终事实来源。双方均不保存原始控件树、截图、PDD 凭据或完整收货地址。
|
||||
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Delivery-Issues
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Delivery-Issues.-
|
||||
wiki_revision: f9245c856a45c10dfbae7fe32b365ca11b008782
|
||||
wiki_revision: 5da7f2250f851caab6d0d1713f869b619ed55eb2
|
||||
synchronized_at: 2026-08-20T08:47:11Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
@@ -53,8 +53,8 @@ synchronized_at: 2026-08-20T08:47:11Z
|
||||
|---|---|---|---|
|
||||
| T29 | [#31](https://git.ilapage.cn/OPC/goauto/issues/31) | PDD 商品采购档案与规格 JSON 管理(2026-08-17 已验收) | 已完成 |
|
||||
| T30 | [#32](https://git.ilapage.cn/OPC/goauto/issues/32) | 采购闭环数据关系与交互原型(2026-08-18 已验收) | 已完成;原型已确认,是全部采购代码依据 |
|
||||
| T31 | [#33](https://git.ilapage.cn/OPC/goauto/issues/33) | 采购任务数据模型与共享 API 契约 | 已实施,等待验收;MySQL 8.4 迁移已验证 |
|
||||
| T32 | [#34](https://git.ilapage.cn/OPC/goauto/issues/34) | 服务端采购任务、租约、幂等与状态机 | #33 |
|
||||
| T31 | [#33](https://git.ilapage.cn/OPC/goauto/issues/33) | 采购任务数据模型与共享 API 契约(2026-08-20 已验收) | 已完成 |
|
||||
| T32 | [#34](https://git.ilapage.cn/OPC/goauto/issues/34) | 服务端采购任务、租约、幂等与状态机 | 已实施,等待验收;未执行真实下单 |
|
||||
| T33 | [#35](https://git.ilapage.cn/OPC/goauto/issues/35) | Admin 采购任务与人工处理页面 | #33、#34 |
|
||||
| T34 | [#36](https://git.ilapage.cn/OPC/goauto/issues/36) | Android 地址后缀、不可逆门禁与创建订单 | #33、#34、#42;真机前再次人工确认 |
|
||||
| T35 | [#37](https://git.ilapage.cn/OPC/goauto/issues/37) | 服务端物流调度与货运宝自动回填 | 采购任务与有效订单能力 |
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"github.com/go-admin-team/go-admin-core/sdk"
|
||||
goautodevice "go-admin/app/goauto/device"
|
||||
goautoproduct "go-admin/app/goauto/product"
|
||||
goautopurchase "go-admin/app/goauto/purchase"
|
||||
goautorule "go-admin/app/goauto/rule"
|
||||
goautoshopeeproduct "go-admin/app/goauto/shopeeproduct"
|
||||
goautosybimport "go-admin/app/goauto/sybimport"
|
||||
@@ -48,6 +49,7 @@ func InitRouter() {
|
||||
// 注册 GoAuto Agent 与设备管理路由。
|
||||
goautodevice.InitRouter(r, authMiddleware)
|
||||
goautotask.InitRouter(r, authMiddleware)
|
||||
goautopurchase.InitRouter(r, authMiddleware)
|
||||
goautoproduct.InitRouter(r, authMiddleware)
|
||||
goautorule.InitRouter(r, authMiddleware)
|
||||
goautoshopeeproduct.InitRouter(r, authMiddleware)
|
||||
|
||||
@@ -47,7 +47,7 @@ func seedPurchaseFixtures(t *testing.T) purchaseFixtures {
|
||||
func newPurchaseTask(fixtures purchaseFixtures, requestID string, status string) models.PurchaseTask {
|
||||
deviceID, accountID, sybID := fixtures.device.ID, fixtures.account.ID, fixtures.syb.ID
|
||||
return models.PurchaseTask{
|
||||
SYBProductID: &sybID, ShopeeProductID: fixtures.shopee.ID, PDDProductID: fixtures.pdd.ID,
|
||||
SYBProductID: &sybID, ShopeeProductID: &fixtures.shopee.ID, PDDProductID: fixtures.pdd.ID,
|
||||
DeviceID: &deviceID, PDDAccountID: &accountID, ExecutionMode: models.PurchaseExecutionModeLive, Status: status,
|
||||
ShopeeItemIDSnapshot: fixtures.shopee.ShopeeItemID, ShopeeTitleSnapshot: fixtures.shopee.Title,
|
||||
ShopeeShopNameSnapshot: fixtures.shopee.ShopName, PDDURLSnapshot: fixtures.pdd.URL,
|
||||
@@ -174,7 +174,8 @@ func TestPurchaseAttemptIdempotency(t *testing.T) {
|
||||
if err := db.Create(&task).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
attempt := models.PurchaseTaskAttempt{TaskID: task.ID, AttemptID: "10000000-0000-0000-0000-000000000001", AttemptNumber: 1, Phase: models.PurchaseAttemptPhasePurchase, Status: models.PurchaseAttemptStatusPending, RuleSnapshotHash: strings.Repeat("a", 64), SpecDecisionSnapshot: `{}`}
|
||||
startRequestID := "20000000-0000-0000-0000-000000000001"
|
||||
attempt := models.PurchaseTaskAttempt{TaskID: task.ID, AttemptID: "10000000-0000-0000-0000-000000000001", AttemptNumber: 1, Phase: models.PurchaseAttemptPhasePurchase, Status: models.PurchaseAttemptStatusPending, RuleSnapshotHash: strings.Repeat("a", 64), SpecDecisionSnapshot: `{}`, StartRequestID: &startRequestID}
|
||||
if err := db.Create(&attempt).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@@ -68,16 +68,16 @@ type PurchaseTask struct {
|
||||
// MySQL 8.4 rejects a column used by both a foreign key with referential
|
||||
// actions and a cross-field CHECK. The live-mode requirement is therefore
|
||||
// enforced by the model hook and the #34 creation service, not a DB CHECK.
|
||||
SYBProductID *uint64 `json:"sybProductId" gorm:"index;uniqueIndex:ux_purchase_task_active_syb,priority:1"`
|
||||
SYBProduct *SYBProduct `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
ShopeeProductID uint64 `json:"shopeeProductId" gorm:"not null;index"`
|
||||
ShopeeProduct ShopeeProduct `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
PDDProductID uint64 `json:"pddProductId" gorm:"not null;index"`
|
||||
PDDProduct PDDProduct `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
DeviceID *uint64 `json:"deviceId" gorm:"index;uniqueIndex:ux_purchase_task_running_device,priority:1"`
|
||||
Device *AgentDevice `json:"-"`
|
||||
PDDAccountID *uint64 `json:"pddAccountId" gorm:"index;uniqueIndex:ux_purchase_task_running_account,priority:1"`
|
||||
PDDAccount *PDDAccount `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
SYBProductID *uint64 `json:"sybProductId" gorm:"index;uniqueIndex:ux_purchase_task_active_syb,priority:1"`
|
||||
SYBProduct *SYBProduct `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
ShopeeProductID *uint64 `json:"shopeeProductId" gorm:"index"`
|
||||
ShopeeProduct *ShopeeProduct `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
PDDProductID uint64 `json:"pddProductId" gorm:"not null;index"`
|
||||
PDDProduct PDDProduct `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
DeviceID *uint64 `json:"deviceId" gorm:"index;uniqueIndex:ux_purchase_task_running_device,priority:1"`
|
||||
Device *AgentDevice `json:"-"`
|
||||
PDDAccountID *uint64 `json:"pddAccountId" gorm:"index;uniqueIndex:ux_purchase_task_running_account,priority:1"`
|
||||
PDDAccount *PDDAccount `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
|
||||
ExecutionMode string `json:"executionMode" gorm:"size:16;not null;check:ck_purchase_task_execution_mode,execution_mode IN ('rehearsal','live')"`
|
||||
Status string `json:"status" gorm:"size:32;not null;index;check:ck_purchase_task_status,status IN ('pending','spec_probe_pending','running','rehearsal_completed','order_submit_started','order_created','order_result_unknown','failed','cancelled')"`
|
||||
@@ -88,18 +88,20 @@ type PurchaseTask struct {
|
||||
DeviceRunSlot *uint8 `json:"-" gorm:"uniqueIndex:ux_purchase_task_running_device,priority:2"`
|
||||
AccountRunSlot *uint8 `json:"-" gorm:"uniqueIndex:ux_purchase_task_running_account,priority:2"`
|
||||
|
||||
ShopeeItemIDSnapshot string `json:"shopeeItemIdSnapshot" gorm:"size:64;not null"`
|
||||
ShopeeTitleSnapshot string `json:"shopeeTitleSnapshot" gorm:"size:500;not null;default:''"`
|
||||
ShopeeShopNameSnapshot string `json:"shopeeShopNameSnapshot" gorm:"size:255;not null;default:''"`
|
||||
PDDURLSnapshot string `json:"pddUrlSnapshot" gorm:"type:text;not null"`
|
||||
PDDGoodsIDSnapshot string `json:"pddGoodsIdSnapshot" gorm:"size:32;not null;index"`
|
||||
PDDTitleSnapshot string `json:"pddTitleSnapshot" gorm:"size:500;not null;default:''"`
|
||||
TargetColorSnapshot string `json:"targetColorSnapshot" gorm:"size:255;not null;default:''"`
|
||||
TargetSizeSnapshot string `json:"targetSizeSnapshot" gorm:"size:255;not null;default:''"`
|
||||
MappedColorSnapshot string `json:"mappedColorSnapshot" gorm:"size:255;not null;default:''"`
|
||||
MappedSizeSnapshot string `json:"mappedSizeSnapshot" gorm:"size:255;not null;default:''"`
|
||||
SpecSource string `json:"specSource" gorm:"size:32;not null;default:unresolved;check:ck_purchase_task_spec_source,spec_source IN ('unresolved','manual_mapping','exact_match','ai_match')"`
|
||||
SpecDecisionSnapshot string `json:"-" gorm:"type:json;not null"`
|
||||
ShopeeItemIDSnapshot string `json:"shopeeItemIdSnapshot" gorm:"size:64;not null"`
|
||||
ShopeeTitleSnapshot string `json:"shopeeTitleSnapshot" gorm:"size:500;not null;default:''"`
|
||||
ShopeeShopNameSnapshot string `json:"shopeeShopNameSnapshot" gorm:"size:255;not null;default:''"`
|
||||
PDDURLSnapshot string `json:"pddUrlSnapshot" gorm:"type:text;not null"`
|
||||
PDDGoodsIDSnapshot string `json:"pddGoodsIdSnapshot" gorm:"size:32;not null;index"`
|
||||
PDDTitleSnapshot string `json:"pddTitleSnapshot" gorm:"size:500;not null;default:''"`
|
||||
TargetColorSnapshot string `json:"targetColorSnapshot" gorm:"size:255;not null;default:''"`
|
||||
TargetSizeSnapshot string `json:"targetSizeSnapshot" gorm:"size:255;not null;default:''"`
|
||||
MappedColorSnapshot string `json:"mappedColorSnapshot" gorm:"size:255;not null;default:''"`
|
||||
MappedSizeSnapshot string `json:"mappedSizeSnapshot" gorm:"size:255;not null;default:''"`
|
||||
SpecSource string `json:"specSource" gorm:"size:32;not null;default:unresolved;check:ck_purchase_task_spec_source,spec_source IN ('unresolved','manual_mapping','exact_match','ai_match')"`
|
||||
SpecDecisionSnapshot string `json:"-" gorm:"type:json;not null"`
|
||||
SpecDecisionRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_spec_decision_request_id"`
|
||||
SpecDecisionBy *uint64 `json:"specDecisionBy"`
|
||||
|
||||
Quantity int64 `json:"quantity" gorm:"not null;check:ck_purchase_task_quantity,quantity >= 1"`
|
||||
ReferenceUnitPriceCent int64 `json:"referenceUnitPriceCent" gorm:"not null;check:ck_purchase_task_reference_price,reference_unit_price_cent >= 0"`
|
||||
@@ -114,10 +116,13 @@ type PurchaseTask struct {
|
||||
PDDAccountRefSnapshot string `json:"pddAccountRefSnapshot" gorm:"size:100;not null;default:''"`
|
||||
AddressSuffix string `json:"addressSuffix" gorm:"size:32;not null;default:''"`
|
||||
|
||||
LeaseExpiresAt *time.Time `json:"leaseExpiresAt" gorm:"index"`
|
||||
LeaseVersion uint64 `json:"leaseVersion" gorm:"not null;default:0"`
|
||||
CreateRequestID string `json:"-" gorm:"size:36;not null;uniqueIndex:ux_purchase_task_create_request_id"`
|
||||
ClaimRequestID *string `json:"-" gorm:"size:36;uniqueIndex:ux_purchase_task_claim_request_id"`
|
||||
LeaseExpiresAt *time.Time `json:"leaseExpiresAt" gorm:"index"`
|
||||
LeaseVersion uint64 `json:"leaseVersion" gorm:"not null;default:0"`
|
||||
CreateRequestID string `json:"-" gorm:"size:36;not null;uniqueIndex:ux_purchase_task_create_request_id"`
|
||||
ClaimRequestID *string `json:"-" gorm:"size:36;uniqueIndex:ux_purchase_task_claim_request_id"`
|
||||
OrderSubmitRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_order_submit_request_id"`
|
||||
StatusVersion uint64 `json:"statusVersion" gorm:"not null;default:1"`
|
||||
StatusChangedAt time.Time `json:"statusChangedAt" gorm:"not null"`
|
||||
|
||||
PDDOrderNo *string `json:"pddOrderNo" gorm:"size:100;index"`
|
||||
OrderSubmittedAt *time.Time `json:"orderSubmittedAt"`
|
||||
@@ -133,18 +138,34 @@ type PurchaseTask struct {
|
||||
WritebackRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_writeback_request_id"`
|
||||
WritebackAt *time.Time `json:"writebackAt"`
|
||||
|
||||
RePurchaseAuthorizedAt *time.Time `json:"rePurchaseAuthorizedAt"`
|
||||
RePurchaseAuthorizedBy *uint64 `json:"rePurchaseAuthorizedBy"`
|
||||
RePurchaseConsumedAt *time.Time `json:"rePurchaseConsumedAt"`
|
||||
ErrorCode *string `json:"errorCode" gorm:"size:64;index"`
|
||||
ErrorMessage *string `json:"errorMessage" gorm:"size:1000"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
RePurchaseAuthorizedAt *time.Time `json:"rePurchaseAuthorizedAt"`
|
||||
RePurchaseAuthorizedBy *uint64 `json:"rePurchaseAuthorizedBy"`
|
||||
RePurchaseConsumedAt *time.Time `json:"rePurchaseConsumedAt"`
|
||||
RePurchaseRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_repurchase_request_id"`
|
||||
PaymentReviewRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_payment_review_request_id"`
|
||||
WritebackSelectRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_writeback_select_request_id"`
|
||||
CancelRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_cancel_request_id"`
|
||||
UnknownResolveRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_task_unknown_resolve_request_id"`
|
||||
CancelledAt *time.Time `json:"cancelledAt"`
|
||||
CancelledBy *uint64 `json:"cancelledBy"`
|
||||
CancelReason *string `json:"cancelReason" gorm:"size:500"`
|
||||
ErrorCode *string `json:"errorCode" gorm:"size:64;index"`
|
||||
ErrorMessage *string `json:"errorMessage" gorm:"size:1000"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
}
|
||||
|
||||
func (PurchaseTask) TableName() string { return "purchase_task" }
|
||||
|
||||
func (task *PurchaseTask) BeforeCreate(_ *gorm.DB) error { return task.syncPurchaseGuardSlots() }
|
||||
func (task *PurchaseTask) BeforeCreate(_ *gorm.DB) error {
|
||||
if task.StatusVersion == 0 {
|
||||
task.StatusVersion = 1
|
||||
}
|
||||
if task.StatusChangedAt.IsZero() {
|
||||
task.StatusChangedAt = time.Now().UTC()
|
||||
}
|
||||
return task.syncPurchaseGuardSlots()
|
||||
}
|
||||
|
||||
func (task *PurchaseTask) BeforeSave(_ *gorm.DB) error { return task.syncPurchaseGuardSlots() }
|
||||
|
||||
@@ -162,7 +183,12 @@ func (task *PurchaseTask) syncPurchaseGuardSlots() error {
|
||||
active = true
|
||||
case PurchaseTaskStatusRunning, PurchaseTaskStatusOrderSubmitStarted:
|
||||
active, running = true, true
|
||||
case PurchaseTaskStatusOrderCreated, PurchaseTaskStatusOrderResultUnknown:
|
||||
case PurchaseTaskStatusOrderCreated:
|
||||
// A successful order occupies the SYB slot until a one-time re-purchase
|
||||
// authorization is consumed. After consumption the old order remains a
|
||||
// historical fact and future edits must not reclaim the new task's slot.
|
||||
active = task.RePurchaseConsumedAt == nil
|
||||
case PurchaseTaskStatusOrderResultUnknown:
|
||||
active = true
|
||||
case PurchaseTaskStatusRehearsalCompleted, PurchaseTaskStatusFailed, PurchaseTaskStatusCancelled:
|
||||
case "":
|
||||
@@ -209,23 +235,26 @@ func (task *PurchaseTask) syncPurchaseGuardSlots() error {
|
||||
// row remains the final business fact. No accessibility tree or screenshot is
|
||||
// stored here.
|
||||
type PurchaseTaskAttempt struct {
|
||||
ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"`
|
||||
TaskID uint64 `json:"taskId" gorm:"not null;index;uniqueIndex:ux_purchase_attempt_number,priority:1"`
|
||||
Task PurchaseTask `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
AttemptID string `json:"attemptId" gorm:"size:36;not null;uniqueIndex:ux_purchase_attempt_id"`
|
||||
AttemptNumber int `json:"attemptNumber" gorm:"not null;uniqueIndex:ux_purchase_attempt_number,priority:2;check:ck_purchase_attempt_number,attempt_number >= 1"`
|
||||
Phase string `json:"phase" gorm:"size:16;not null;check:ck_purchase_attempt_phase,phase IN ('spec_probe','purchase')"`
|
||||
Status string `json:"status" gorm:"size:16;not null;index;check:ck_purchase_attempt_status,status IN ('pending','running','completed','failed')"`
|
||||
DeviceID *uint64 `json:"deviceId" gorm:"index"`
|
||||
RuleSnapshotHash string `json:"ruleSnapshotHash" gorm:"size:64;not null"`
|
||||
SpecDecisionSnapshot string `json:"-" gorm:"type:json;not null"`
|
||||
ResultRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_attempt_result_request_id"`
|
||||
ErrorCode *string `json:"errorCode" gorm:"size:64;index"`
|
||||
ErrorMessage *string `json:"errorMessage" gorm:"size:1000"`
|
||||
StartedAt *time.Time `json:"startedAt"`
|
||||
FinishedAt *time.Time `json:"finishedAt"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"`
|
||||
TaskID uint64 `json:"taskId" gorm:"not null;index;uniqueIndex:ux_purchase_attempt_number,priority:1"`
|
||||
Task *PurchaseTask `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"`
|
||||
AttemptID string `json:"attemptId" gorm:"size:36;not null;uniqueIndex:ux_purchase_attempt_id"`
|
||||
AttemptNumber int `json:"attemptNumber" gorm:"not null;uniqueIndex:ux_purchase_attempt_number,priority:2;check:ck_purchase_attempt_number,attempt_number >= 1"`
|
||||
Phase string `json:"phase" gorm:"size:16;not null;check:ck_purchase_attempt_phase,phase IN ('spec_probe','purchase')"`
|
||||
Status string `json:"status" gorm:"size:16;not null;index;check:ck_purchase_attempt_status,status IN ('pending','running','completed','failed')"`
|
||||
DeviceID *uint64 `json:"deviceId" gorm:"index"`
|
||||
RuleSnapshotHash string `json:"ruleSnapshotHash" gorm:"size:64;not null"`
|
||||
SpecDecisionSnapshot string `json:"-" gorm:"type:json;not null"`
|
||||
ResultRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_attempt_result_request_id"`
|
||||
StartRequestID *string `json:"-" gorm:"size:64;uniqueIndex:ux_purchase_attempt_start_request_id"`
|
||||
ResultHash *string `json:"-" gorm:"size:64"`
|
||||
ResultType *string `json:"-" gorm:"size:32"`
|
||||
ErrorCode *string `json:"errorCode" gorm:"size:64;index"`
|
||||
ErrorMessage *string `json:"errorMessage" gorm:"size:1000"`
|
||||
StartedAt *time.Time `json:"startedAt"`
|
||||
FinishedAt *time.Time `json:"finishedAt"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
}
|
||||
|
||||
func (PurchaseTaskAttempt) TableName() string { return "purchase_task_attempt" }
|
||||
|
||||
@@ -0,0 +1,250 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/app/goauto/models"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/go-admin-team/go-admin-core/sdk/pkg"
|
||||
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type Handler struct{ DB *gorm.DB }
|
||||
|
||||
func (h Handler) AdminCreate(c *gin.Context) {
|
||||
if !allowedOperator(c) {
|
||||
return
|
||||
}
|
||||
var req CreateRequest
|
||||
if !decode(c, &req) {
|
||||
return
|
||||
}
|
||||
service, ok := h.service(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
record, replayed, err := service.Create(c.Request.Context(), req)
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"data": record, "replayed": replayed})
|
||||
}
|
||||
func (h Handler) Next(c *gin.Context) {
|
||||
service, ok := h.service(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
p, err := service.Next(c.Request.Context(), bearer(c.GetHeader("Authorization")))
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
if p == nil {
|
||||
c.Status(http.StatusNoContent)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"data": p})
|
||||
}
|
||||
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) }
|
||||
func (h Handler) Result(c *gin.Context) {
|
||||
id, ok := pathID(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var req ResultRequest
|
||||
if !decode(c, &req) {
|
||||
return
|
||||
}
|
||||
service, serviceOK := h.service(c)
|
||||
if !serviceOK {
|
||||
return
|
||||
}
|
||||
p, err := service.SubmitResult(c.Request.Context(), id, req, bearer(c.GetHeader("Authorization")))
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"data": p})
|
||||
}
|
||||
func (h Handler) SpecDecision(c *gin.Context) {
|
||||
if !allowedOperator(c) {
|
||||
return
|
||||
}
|
||||
id, ok := pathID(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var req SpecDecisionRequest
|
||||
if !decode(c, &req) {
|
||||
return
|
||||
}
|
||||
req.OperatorID = operatorID(c)
|
||||
service, serviceOK := h.service(c)
|
||||
if !serviceOK {
|
||||
return
|
||||
}
|
||||
p, replayed, err := service.ApplySpecDecision(c.Request.Context(), id, req)
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"data": p, "replayed": replayed})
|
||||
}
|
||||
func (h Handler) AuthorizeRePurchase(c *gin.Context) { h.manual(c, (*Service).AuthorizeRePurchase) }
|
||||
func (h Handler) ReviewPayment(c *gin.Context) { h.manual(c, (*Service).ReviewPayment) }
|
||||
func (h Handler) SelectWriteback(c *gin.Context) { h.manual(c, (*Service).SelectWriteback) }
|
||||
func (h Handler) Cancel(c *gin.Context) { h.manual(c, (*Service).Cancel) }
|
||||
func (h Handler) ResolveUnknown(c *gin.Context) { h.manual(c, (*Service).ResolveUnknown) }
|
||||
|
||||
func (h Handler) action(c *gin.Context, fn func(*Service, context.Context, uint64, ActionRequest, string) (TaskPayload, error)) {
|
||||
id, ok := pathID(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var req ActionRequest
|
||||
if !decode(c, &req) {
|
||||
return
|
||||
}
|
||||
service, serviceOK := h.service(c)
|
||||
if !serviceOK {
|
||||
return
|
||||
}
|
||||
p, err := fn(service, c.Request.Context(), id, req, bearer(c.GetHeader("Authorization")))
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"data": p})
|
||||
}
|
||||
|
||||
func (h Handler) manual(c *gin.Context, fn func(*Service, context.Context, uint64, ManualRequest) (models.PurchaseTask, bool, error)) {
|
||||
if !allowedOperator(c) {
|
||||
return
|
||||
}
|
||||
id, ok := pathID(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var req ManualRequest
|
||||
if !decode(c, &req) {
|
||||
return
|
||||
}
|
||||
req.OperatorID = operatorID(c)
|
||||
if req.OperatorID == 0 {
|
||||
writeError(c, fail(CodeInvalidRequest, "无法识别当前操作人"))
|
||||
return
|
||||
}
|
||||
service, serviceOK := h.service(c)
|
||||
if !serviceOK {
|
||||
return
|
||||
}
|
||||
p, replayed, err := fn(service, c.Request.Context(), id, req)
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"data": p, "replayed": replayed})
|
||||
}
|
||||
|
||||
func allowedOperator(c *gin.Context) bool {
|
||||
role, _ := jwt.ExtractClaims(c)["rolekey"].(string)
|
||||
if role == "admin" || role == "purchaser" {
|
||||
return true
|
||||
}
|
||||
c.JSON(http.StatusForbidden, gin.H{"code": "FORBIDDEN", "message": "只有管理员或采购员可以操作采购任务"})
|
||||
c.Abort()
|
||||
return false
|
||||
}
|
||||
|
||||
func (h Handler) service(c *gin.Context) (*Service, bool) {
|
||||
if h.DB != nil {
|
||||
return NewService(h.DB), true
|
||||
}
|
||||
db, err := pkg.GetOrm(c)
|
||||
if err != nil {
|
||||
writeError(c, internal(err))
|
||||
return nil, false
|
||||
}
|
||||
return NewService(db), true
|
||||
}
|
||||
|
||||
func decode(c *gin.Context, v any) bool {
|
||||
c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, 1<<20)
|
||||
d := json.NewDecoder(c.Request.Body)
|
||||
d.DisallowUnknownFields()
|
||||
if err := d.Decode(v); err != nil {
|
||||
writeError(c, fail(CodeInvalidRequest, "请求 JSON 无效"))
|
||||
return false
|
||||
}
|
||||
if err := d.Decode(&struct{}{}); !errors.Is(err, io.EOF) {
|
||||
writeError(c, fail(CodeInvalidRequest, "请求只能包含一个 JSON 对象"))
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
func pathID(c *gin.Context) (uint64, bool) {
|
||||
id, err := strconv.ParseUint(c.Param("taskId"), 10, 64)
|
||||
if err != nil || id == 0 {
|
||||
writeError(c, fail(CodeInvalidRequest, "taskId 无效"))
|
||||
return 0, false
|
||||
}
|
||||
return id, true
|
||||
}
|
||||
func bearer(v string) string {
|
||||
p := strings.Fields(v)
|
||||
if len(p) == 2 && strings.EqualFold(p[0], "Bearer") {
|
||||
return p[1]
|
||||
}
|
||||
return ""
|
||||
}
|
||||
func writeError(c *gin.Context, err error) {
|
||||
code, msg, retryable, status := CodeInternal, "服务端处理失败", true, http.StatusInternalServerError
|
||||
var se *ServiceError
|
||||
var de *device.ServiceError
|
||||
if errors.As(err, &se) {
|
||||
code, msg, retryable = se.Code, se.Message, se.Retryable
|
||||
} else if errors.As(err, &de) {
|
||||
code, msg, retryable = de.Code, de.Message, de.Retryable
|
||||
}
|
||||
switch code {
|
||||
case CodeInvalidRequest:
|
||||
status = http.StatusUnprocessableEntity
|
||||
case device.CodeTokenInvalid:
|
||||
status = http.StatusUnauthorized
|
||||
case device.CodeDeviceDisabled:
|
||||
status = http.StatusForbidden
|
||||
case CodeTaskNotFound:
|
||||
status = http.StatusNotFound
|
||||
case CodeStateConflict, CodeCapabilityMismatch, CodeDeviceBusy, CodeTaskClaimed, CodeLeaseExpired, CodeMappingRequired, CodeResultConflict, CodeRePurchaseRequired:
|
||||
status = http.StatusConflict
|
||||
}
|
||||
c.JSON(status, gin.H{"code": code, "message": msg, "retryable": retryable})
|
||||
}
|
||||
func operatorID(c *gin.Context) uint64 {
|
||||
claims := jwt.ExtractClaims(c)
|
||||
switch v := claims["identity"].(type) {
|
||||
case float64:
|
||||
return uint64(v)
|
||||
case int:
|
||||
return uint64(v)
|
||||
case json.Number:
|
||||
n, _ := strconv.ParseUint(string(v), 10, 64)
|
||||
return n
|
||||
case string:
|
||||
n, _ := strconv.ParseUint(v, 10, 64)
|
||||
return n
|
||||
}
|
||||
return 0
|
||||
}
|
||||
@@ -0,0 +1,488 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/app/goauto/models"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
func (s *Service) Next(ctx context.Context, token string) (*TaskPayload, error) {
|
||||
d, err := device.NewService(s.DB).Authenticate(ctx, token)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if d.Status != models.DeviceStatusOnline {
|
||||
return nil, fail(CodeStateConflict, "设备离线,不能领取采购任务")
|
||||
}
|
||||
var running models.PurchaseTask
|
||||
if err = s.DB.WithContext(ctx).Where("device_id = ? AND status IN ?", d.ID, []string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted}).First(&running).Error; err == nil {
|
||||
return s.payload(running, nil, false)
|
||||
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, internal(err)
|
||||
}
|
||||
now := s.Now()
|
||||
var candidates []models.PurchaseTask
|
||||
if err = s.DB.WithContext(ctx).Where("status IN ? AND (lease_expires_at IS NULL OR lease_expires_at <= ?) AND (device_id IS NULL OR device_id = ?)", []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending}, now, d.ID).Order("CASE WHEN device_id IS NULL THEN 1 ELSE 0 END, created_at, id").Limit(100).Find(&candidates).Error; err != nil {
|
||||
return nil, internal(err)
|
||||
}
|
||||
for _, t := range candidates {
|
||||
if t.Status == models.PurchaseTaskStatusSpecProbePending && t.MappedColorSnapshot == "" && t.MappedSizeSnapshot == "" {
|
||||
continue
|
||||
}
|
||||
required, er := decodeStrings(t.RequiredCapabilitiesJSON)
|
||||
if er != nil {
|
||||
return nil, internal(er)
|
||||
}
|
||||
if er = ensureCapabilities(d, required); er == nil {
|
||||
return s.payload(t, nil, false)
|
||||
}
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (s *Service) Claim(ctx context.Context, taskID uint64, req ActionRequest, token string) (TaskPayload, error) {
|
||||
if _, err := uuid.Parse(req.RequestID); err != nil {
|
||||
return TaskPayload{}, fail(CodeInvalidRequest, "requestId 无效")
|
||||
}
|
||||
var out TaskPayload
|
||||
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
d, err := device.NewService(tx).Authenticate(ctx, token)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if d.Status != models.DeviceStatusOnline {
|
||||
return fail(CodeStateConflict, "设备离线,不能领取采购任务")
|
||||
}
|
||||
var t models.PurchaseTask
|
||||
if err = tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&t, taskID).Error; err != nil {
|
||||
return purchaseNotFound(err)
|
||||
}
|
||||
if t.ClaimRequestID != nil && *t.ClaimRequestID == req.RequestID && t.DeviceID != nil && *t.DeviceID == d.ID {
|
||||
p, e := s.payload(t, nil, true)
|
||||
out = *p
|
||||
return e
|
||||
}
|
||||
if t.Status != models.PurchaseTaskStatusPending && t.Status != models.PurchaseTaskStatusSpecProbePending {
|
||||
return fail(CodeStateConflict, "任务当前状态不能领取")
|
||||
}
|
||||
if t.DeviceID != nil && *t.DeviceID != d.ID {
|
||||
return fail(CodeStateConflict, "任务已指定给其他设备")
|
||||
}
|
||||
if t.LeaseExpiresAt != nil && t.LeaseExpiresAt.After(s.Now()) {
|
||||
return fail(CodeTaskClaimed, "任务已被领取")
|
||||
}
|
||||
if err = tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("status IN ? AND lease_expires_at <= ?", []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending}, s.Now()).Updates(map[string]any{"device_run_slot": nil, "account_run_slot": nil}).Error; err != nil {
|
||||
return internal(err)
|
||||
}
|
||||
required, e := decodeStrings(t.RequiredCapabilitiesJSON)
|
||||
if e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
if e = ensureCapabilities(d, required); e != nil {
|
||||
return e
|
||||
}
|
||||
if e = ensureDeviceFree(tx, d.ID, t.ID, s.Now()); e != nil {
|
||||
return e
|
||||
}
|
||||
if e = ensureAccountFree(tx, t.PDDAccountID, t.ID, s.Now()); e != nil {
|
||||
return e
|
||||
}
|
||||
lease := s.Now().Add(s.lease())
|
||||
one := uint8(1)
|
||||
updates := map[string]any{"device_id": d.ID, "lease_expires_at": lease, "lease_version": gorm.Expr("lease_version + 1"), "claim_request_id": req.RequestID, "device_run_slot": one}
|
||||
if t.PDDAccountID != nil {
|
||||
updates["account_run_slot"] = one
|
||||
}
|
||||
result := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ? AND status IN ? AND (lease_expires_at IS NULL OR lease_expires_at <= ?)", t.ID, []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending}, s.Now()).Updates(updates)
|
||||
if result.Error != nil {
|
||||
return conflictOrInternal(result.Error)
|
||||
}
|
||||
if result.RowsAffected != 1 {
|
||||
return fail(CodeTaskClaimed, "任务已被其他设备领取")
|
||||
}
|
||||
if e = tx.First(&t, t.ID).Error; e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
p, e := s.payload(t, nil, false)
|
||||
out = *p
|
||||
return e
|
||||
})
|
||||
return out, err
|
||||
}
|
||||
|
||||
func (s *Service) Start(ctx context.Context, taskID uint64, req ActionRequest, token string) (TaskPayload, error) {
|
||||
if _, err := uuid.Parse(req.RequestID); err != nil {
|
||||
return TaskPayload{}, fail(CodeInvalidRequest, "requestId 无效")
|
||||
}
|
||||
var out TaskPayload
|
||||
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
d, err := device.NewService(tx).Authenticate(ctx, token)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var t models.PurchaseTask
|
||||
if err = tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&t, taskID).Error; err != nil {
|
||||
return purchaseNotFound(err)
|
||||
}
|
||||
var existing models.PurchaseTaskAttempt
|
||||
if err = tx.Where("start_request_id = ?", req.RequestID).First(&existing).Error; err == nil {
|
||||
if existing.TaskID != t.ID || existing.DeviceID == nil || *existing.DeviceID != d.ID {
|
||||
return fail(CodeResultConflict, "start requestId 已用于其他任务")
|
||||
}
|
||||
p, e := s.payload(t, &existing, true)
|
||||
out = *p
|
||||
return e
|
||||
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return internal(err)
|
||||
}
|
||||
if t.Status != models.PurchaseTaskStatusPending && t.Status != models.PurchaseTaskStatusSpecProbePending {
|
||||
return fail(CodeStateConflict, "任务当前状态不能开始")
|
||||
}
|
||||
if t.DeviceID == nil || *t.DeviceID != d.ID {
|
||||
return fail(CodeStateConflict, "任务不属于当前设备")
|
||||
}
|
||||
if t.LeaseExpiresAt == nil || !t.LeaseExpiresAt.After(s.Now()) {
|
||||
return fail(CodeLeaseExpired, "任务租约已过期,请重新领取")
|
||||
}
|
||||
required, e := decodeStrings(t.RequiredCapabilitiesJSON)
|
||||
if e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
if e = ensureCapabilities(d, required); e != nil {
|
||||
return e
|
||||
}
|
||||
if e = ensureDeviceFree(tx, d.ID, t.ID, s.Now()); e != nil {
|
||||
return e
|
||||
}
|
||||
if e = ensureAccountFree(tx, t.PDDAccountID, t.ID, s.Now()); e != nil {
|
||||
return e
|
||||
}
|
||||
phase := models.PurchaseAttemptPhasePurchase
|
||||
if t.SpecSource == "unresolved" {
|
||||
phase = models.PurchaseAttemptPhaseSpecProbe
|
||||
}
|
||||
var count int64
|
||||
if e = tx.Model(&models.PurchaseTaskAttempt{}).Where("task_id = ?", t.ID).Count(&count).Error; e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
h := sha256.Sum256([]byte(t.RuleSnapshot))
|
||||
now := s.Now()
|
||||
a := models.PurchaseTaskAttempt{TaskID: t.ID, AttemptID: uuid.NewString(), AttemptNumber: int(count) + 1, Phase: phase, Status: models.PurchaseAttemptStatusRunning, DeviceID: &d.ID, RuleSnapshotHash: hex.EncodeToString(h[:]), SpecDecisionSnapshot: t.SpecDecisionSnapshot, StartRequestID: &req.RequestID, StartedAt: &now}
|
||||
if e = tx.Omit("Task").Create(&a).Error; e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
if e = t.SetStatus(models.PurchaseTaskStatusRunning); e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
t.StatusVersion++
|
||||
t.StatusChangedAt = now
|
||||
t.LeaseExpiresAt = ptrTime(now.Add(s.lease()))
|
||||
if e = tx.Save(&t).Error; e != nil {
|
||||
return conflictOrInternal(e)
|
||||
}
|
||||
p, e := s.payload(t, &a, false)
|
||||
out = *p
|
||||
return e
|
||||
})
|
||||
return out, err
|
||||
}
|
||||
|
||||
func (s *Service) MarkOrderSubmitStarted(ctx context.Context, taskID uint64, req ActionRequest, token string) (TaskPayload, error) {
|
||||
return s.withRunning(ctx, taskID, token, func(tx *gorm.DB, t *models.PurchaseTask, a *models.PurchaseTaskAttempt, d models.AgentDevice) (TaskPayload, error) {
|
||||
if t.ExecutionMode != models.PurchaseExecutionModeLive {
|
||||
return TaskPayload{}, fail(CodeStateConflict, "演练任务不能创建订单")
|
||||
}
|
||||
if _, e := uuid.Parse(req.RequestID); e != nil {
|
||||
return TaskPayload{}, fail(CodeInvalidRequest, "requestId 无效")
|
||||
}
|
||||
if t.OrderSubmitRequestID != nil {
|
||||
if *t.OrderSubmitRequestID == req.RequestID {
|
||||
return valuePayload(s, t, a, true)
|
||||
}
|
||||
return TaskPayload{}, fail(CodeResultConflict, "订单提交状态已记录")
|
||||
}
|
||||
now := s.Now()
|
||||
if e := t.SetStatus(models.PurchaseTaskStatusOrderSubmitStarted); e != nil {
|
||||
return TaskPayload{}, internal(e)
|
||||
}
|
||||
t.OrderSubmitRequestID = &req.RequestID
|
||||
t.IrreversibleAt = &now
|
||||
t.StatusVersion++
|
||||
t.StatusChangedAt = now
|
||||
if e := tx.Save(t).Error; e != nil {
|
||||
return TaskPayload{}, conflictOrInternal(e)
|
||||
}
|
||||
return valuePayload(s, t, a, false)
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Service) SubmitResult(ctx context.Context, taskID uint64, req ResultRequest, token string) (TaskPayload, error) {
|
||||
raw, _ := json.Marshal(req)
|
||||
digest := hashBytes(raw)
|
||||
deviceRecord, authErr := device.NewService(s.DB).Authenticate(ctx, token)
|
||||
if authErr != nil {
|
||||
return TaskPayload{}, authErr
|
||||
}
|
||||
var replayTask models.PurchaseTask
|
||||
var replayAttempt models.PurchaseTaskAttempt
|
||||
if err := s.DB.WithContext(ctx).First(&replayTask, taskID).Error; err == nil {
|
||||
if replayTask.DeviceID == nil || *replayTask.DeviceID != deviceRecord.ID {
|
||||
return TaskPayload{}, fail(CodeStateConflict, "任务不属于当前设备")
|
||||
}
|
||||
if err := s.DB.WithContext(ctx).Where("task_id = ? AND attempt_id = ?", taskID, req.TaskAttemptID).First(&replayAttempt).Error; err == nil && replayAttempt.ResultRequestID != nil {
|
||||
if *replayAttempt.ResultRequestID == req.RequestID && replayAttempt.ResultHash != nil && *replayAttempt.ResultHash == digest {
|
||||
return valuePayload(s, &replayTask, &replayAttempt, true)
|
||||
}
|
||||
return TaskPayload{}, fail(CodeResultConflict, "同一次执行已提交不同结果")
|
||||
}
|
||||
}
|
||||
return s.withRunning(ctx, taskID, token, func(tx *gorm.DB, t *models.PurchaseTask, a *models.PurchaseTaskAttempt, d models.AgentDevice) (TaskPayload, error) {
|
||||
if req.TaskAttemptID == "" || req.TaskAttemptID != a.AttemptID {
|
||||
return TaskPayload{}, fail(CodeStateConflict, "taskAttemptId 与当前执行不一致")
|
||||
}
|
||||
if strings.TrimSpace(req.RequestID) == "" {
|
||||
return TaskPayload{}, fail(CodeInvalidRequest, "requestId 不能为空")
|
||||
}
|
||||
if a.ResultRequestID != nil {
|
||||
if *a.ResultRequestID == req.RequestID && a.ResultHash != nil && *a.ResultHash == digest {
|
||||
return valuePayload(s, t, a, true)
|
||||
}
|
||||
return TaskPayload{}, fail(CodeResultConflict, "同一次执行已提交不同结果")
|
||||
}
|
||||
now := s.Now()
|
||||
next := ""
|
||||
switch req.ResultType {
|
||||
case "spec_probe_completed":
|
||||
if len(req.ProbedSpecs) == 0 {
|
||||
return TaskPayload{}, fail(CodeInvalidRequest, "规格探测结果无效")
|
||||
}
|
||||
next = models.PurchaseTaskStatusSpecProbePending
|
||||
a.Status = models.PurchaseAttemptStatusCompleted
|
||||
t.SpecSource = "unresolved"
|
||||
t.MappedColorSnapshot = ""
|
||||
t.MappedSizeSnapshot = ""
|
||||
case "rehearsal_completed":
|
||||
if t.ExecutionMode != models.PurchaseExecutionModeRehearsal {
|
||||
return TaskPayload{}, fail(CodeStateConflict, "正式任务不能提交演练结果")
|
||||
}
|
||||
next = models.PurchaseTaskStatusRehearsalCompleted
|
||||
a.Status = models.PurchaseAttemptStatusCompleted
|
||||
case "order_created":
|
||||
if t.ExecutionMode != models.PurchaseExecutionModeLive || t.Status != models.PurchaseTaskStatusOrderSubmitStarted || strings.TrimSpace(req.PDDOrderNo) == "" || req.OrderSubmittedAt == nil {
|
||||
return TaskPayload{}, fail(CodeInvalidRequest, "订单号或下单时间缺失")
|
||||
}
|
||||
next = models.PurchaseTaskStatusOrderCreated
|
||||
a.Status = models.PurchaseAttemptStatusCompleted
|
||||
t.PDDOrderNo = &req.PDDOrderNo
|
||||
t.OrderSubmittedAt = req.OrderSubmittedAt
|
||||
case "order_result_unknown":
|
||||
if t.ExecutionMode != models.PurchaseExecutionModeLive || t.Status != models.PurchaseTaskStatusOrderSubmitStarted {
|
||||
return TaskPayload{}, fail(CodeStateConflict, "当前任务不能标记订单结果未知")
|
||||
}
|
||||
next = models.PurchaseTaskStatusOrderResultUnknown
|
||||
a.Status = models.PurchaseAttemptStatusFailed
|
||||
case "failed":
|
||||
next = models.PurchaseTaskStatusFailed
|
||||
a.Status = models.PurchaseAttemptStatusFailed
|
||||
t.ErrorCode = &req.ErrorCode
|
||||
t.ErrorMessage = &req.ErrorMessage
|
||||
default:
|
||||
return TaskPayload{}, fail(CodeInvalidRequest, "resultType 无效")
|
||||
}
|
||||
if e := t.SetStatus(next); e != nil {
|
||||
return TaskPayload{}, internal(e)
|
||||
}
|
||||
t.LeaseExpiresAt = nil
|
||||
t.StatusVersion++
|
||||
t.StatusChangedAt = now
|
||||
a.ResultRequestID = &req.RequestID
|
||||
a.ResultHash = &digest
|
||||
a.ResultType = &req.ResultType
|
||||
a.FinishedAt = &now
|
||||
if req.ResultType == "spec_probe_completed" {
|
||||
a.SpecDecisionSnapshot = string(req.ProbedSpecs)
|
||||
}
|
||||
if e := tx.Omit("Task").Save(a).Error; e != nil {
|
||||
return TaskPayload{}, internal(e)
|
||||
}
|
||||
if e := tx.Save(t).Error; e != nil {
|
||||
return TaskPayload{}, conflictOrInternal(e)
|
||||
}
|
||||
return valuePayload(s, t, a, false)
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Service) ApplySpecDecision(ctx context.Context, taskID uint64, req SpecDecisionRequest) (models.PurchaseTask, bool, error) {
|
||||
var t models.PurchaseTask
|
||||
replayed := false
|
||||
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if e := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&t, taskID).Error; e != nil {
|
||||
return purchaseNotFound(e)
|
||||
}
|
||||
if strings.TrimSpace(req.RequestID) == "" || req.OperatorID == 0 {
|
||||
return fail(CodeInvalidRequest, "requestId 或操作人无效")
|
||||
}
|
||||
var a models.PurchaseTaskAttempt
|
||||
if e := tx.Where("task_id = ? AND attempt_id = ?", t.ID, req.TaskAttemptID).First(&a).Error; e != nil {
|
||||
return purchaseNotFound(e)
|
||||
}
|
||||
if a.Status != models.PurchaseAttemptStatusCompleted || a.ResultType == nil || *a.ResultType != "spec_probe_completed" {
|
||||
return fail(CodeStateConflict, "该 attempt 没有可固化的规格探测结果")
|
||||
}
|
||||
if t.SpecDecisionRequestID != nil {
|
||||
same := t.MappedColorSnapshot == req.MappedColor && t.MappedSizeSnapshot == req.MappedSize
|
||||
if same && *t.SpecDecisionRequestID == req.RequestID {
|
||||
replayed = true
|
||||
return nil
|
||||
}
|
||||
return fail(CodeResultConflict, "同一次规格探测的决策已经固化,不能修改")
|
||||
}
|
||||
if t.Status != models.PurchaseTaskStatusSpecProbePending {
|
||||
return fail(CodeStateConflict, "任务不在待规格决策状态")
|
||||
}
|
||||
if req.Source != "ai_match" && req.Source != "exact_match" && req.Source != "manual_mapping" {
|
||||
return fail(CodeInvalidRequest, "规格决策来源无效")
|
||||
}
|
||||
decision := req.Decision
|
||||
if len(decision) == 0 {
|
||||
decision = []byte("{}")
|
||||
}
|
||||
t.MappedColorSnapshot = strings.TrimSpace(req.MappedColor)
|
||||
t.MappedSizeSnapshot = strings.TrimSpace(req.MappedSize)
|
||||
if !req.NoMatch && ((t.TargetColorSnapshot != "" && t.MappedColorSnapshot == "") || (t.TargetSizeSnapshot != "" && t.MappedSizeSnapshot == "")) {
|
||||
return fail(CodeMappingRequired, "没有找到可用的商品规格")
|
||||
}
|
||||
t.SpecSource = req.Source
|
||||
t.SpecDecisionSnapshot = string(decision)
|
||||
t.SpecDecisionRequestID = &req.RequestID
|
||||
t.SpecDecisionBy = &req.OperatorID
|
||||
if req.NoMatch {
|
||||
code, message := "PURCHASE_SPEC_NOT_MATCHED", "没有找到可用的商品规格"
|
||||
t.ErrorCode, t.ErrorMessage = &code, &message
|
||||
if e := t.SetStatus(models.PurchaseTaskStatusFailed); e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
}
|
||||
t.StatusVersion++
|
||||
t.StatusChangedAt = s.Now()
|
||||
return tx.Save(&t).Error
|
||||
})
|
||||
return t, replayed, err
|
||||
}
|
||||
|
||||
func (s *Service) withRunning(ctx context.Context, taskID uint64, token string, fn func(*gorm.DB, *models.PurchaseTask, *models.PurchaseTaskAttempt, models.AgentDevice) (TaskPayload, error)) (TaskPayload, error) {
|
||||
var out TaskPayload
|
||||
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
d, e := device.NewService(tx).Authenticate(ctx, token)
|
||||
if e != nil {
|
||||
return e
|
||||
}
|
||||
var t models.PurchaseTask
|
||||
if e = tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&t, taskID).Error; e != nil {
|
||||
return purchaseNotFound(e)
|
||||
}
|
||||
if t.DeviceID == nil || *t.DeviceID != d.ID {
|
||||
return fail(CodeStateConflict, "任务不属于当前设备")
|
||||
}
|
||||
if t.Status != models.PurchaseTaskStatusRunning && t.Status != models.PurchaseTaskStatusOrderSubmitStarted {
|
||||
return fail(CodeStateConflict, "任务当前状态不能提交结果")
|
||||
}
|
||||
if t.LeaseExpiresAt == nil || !t.LeaseExpiresAt.After(s.Now()) {
|
||||
return fail(CodeLeaseExpired, "任务租约已过期")
|
||||
}
|
||||
required, e := decodeStrings(t.RequiredCapabilitiesJSON)
|
||||
if e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
if e = ensureCapabilities(d, required); e != nil {
|
||||
return e
|
||||
}
|
||||
var a models.PurchaseTaskAttempt
|
||||
if e = tx.Where("task_id = ? AND status = ?", t.ID, models.PurchaseAttemptStatusRunning).Order("attempt_number DESC").First(&a).Error; e != nil {
|
||||
return purchaseNotFound(e)
|
||||
}
|
||||
p, e := fn(tx, &t, &a, d)
|
||||
out = p
|
||||
return e
|
||||
})
|
||||
return out, err
|
||||
}
|
||||
|
||||
func ensureDeviceFree(tx *gorm.DB, deviceID, taskID uint64, now time.Time) error {
|
||||
var count int64
|
||||
if e := tx.Model(&models.PurchaseTask{}).Where("id <> ? AND device_id = ? AND (status IN ? OR (status IN ? AND lease_expires_at > ?))", taskID, deviceID, []string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted}, []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending}, now).Count(&count).Error; e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
if count > 0 {
|
||||
return fail(CodeDeviceBusy, "设备已有采购任务")
|
||||
}
|
||||
if e := tx.Model(&models.CollectionTask{}).Where("device_id = ? AND (status = ? OR (status = ? AND lease_expires_at > ?))", deviceID, models.TaskStatusRunning, models.TaskStatusPending, now).Count(&count).Error; e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
if count > 0 {
|
||||
return fail(CodeDeviceBusy, "设备已有采集任务")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func ensureAccountFree(tx *gorm.DB, accountID *uint64, taskID uint64, now time.Time) error {
|
||||
if accountID == nil {
|
||||
return nil
|
||||
}
|
||||
var count int64
|
||||
if e := tx.Model(&models.PurchaseTask{}).Where("id <> ? AND pdd_account_id = ? AND (status IN ? OR (status IN ? AND lease_expires_at > ?))", taskID, *accountID, []string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted}, []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending}, now).Count(&count).Error; e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
if count > 0 {
|
||||
return fail(CodeDeviceBusy, "拼多多账号已有采购任务")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (s *Service) payload(t models.PurchaseTask, a *models.PurchaseTaskAttempt, replayed bool) (*TaskPayload, error) {
|
||||
p := &TaskPayload{TaskID: t.ID, ExecutionMode: t.ExecutionMode, Status: t.Status, DeviceID: t.DeviceID, PDDProductID: t.PDDProductID, PDDURL: t.PDDURLSnapshot, PDDGoodsID: t.PDDGoodsIDSnapshot, TargetColor: t.TargetColorSnapshot, TargetSize: t.TargetSizeSnapshot, MappedColor: t.MappedColorSnapshot, MappedSize: t.MappedSizeSnapshot, Quantity: t.Quantity, MinUnitPriceCent: t.MinUnitPriceCent, MaxUnitPriceCent: t.MaxUnitPriceCent, Currency: t.Currency, AddressSuffix: t.AddressSuffix, RuleSnapshot: json.RawMessage(t.RuleSnapshot), LeaseExpiresAt: t.LeaseExpiresAt, LeaseVersion: t.LeaseVersion, Replayed: replayed}
|
||||
if a != nil {
|
||||
p.TaskAttemptID = a.AttemptID
|
||||
p.AttemptNumber = a.AttemptNumber
|
||||
p.Phase = a.Phase
|
||||
p.RuleSnapshotHash = a.RuleSnapshotHash
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
func valuePayload(s *Service, t *models.PurchaseTask, a *models.PurchaseTaskAttempt, replayed bool) (TaskPayload, error) {
|
||||
p, e := s.payload(*t, a, replayed)
|
||||
return *p, e
|
||||
}
|
||||
func purchaseNotFound(err error) error {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return fail(CodeTaskNotFound, "采购任务不存在")
|
||||
}
|
||||
return internal(err)
|
||||
}
|
||||
func conflictOrInternal(err error) error {
|
||||
if isDuplicate(err) {
|
||||
return fail(CodeDeviceBusy, "设备或拼多多账号已有运行任务")
|
||||
}
|
||||
return internal(err)
|
||||
}
|
||||
func decodeStrings(raw string) ([]string, error) {
|
||||
var v []string
|
||||
err := json.Unmarshal([]byte(raw), &v)
|
||||
return v, err
|
||||
}
|
||||
func ptrTime(v time.Time) *time.Time { return &v }
|
||||
func (s *Service) lease() time.Duration {
|
||||
if s.LeaseDuration <= 0 {
|
||||
return DefaultLeaseDuration
|
||||
}
|
||||
return s.LeaseDuration
|
||||
}
|
||||
@@ -0,0 +1,166 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
|
||||
"go-admin/app/goauto/models"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
func (s *Service) AuthorizeRePurchase(ctx context.Context, id uint64, req ManualRequest) (models.PurchaseTask, bool, error) {
|
||||
return s.manual(ctx, id, req, func(t *models.PurchaseTask) error {
|
||||
if t.Status != models.PurchaseTaskStatusOrderCreated && t.Status != models.PurchaseTaskStatusFailed && t.Status != models.PurchaseTaskStatusCancelled {
|
||||
return fail(CodeStateConflict, "当前任务不能授权重新采购")
|
||||
}
|
||||
if t.RePurchaseRequestID != nil {
|
||||
return fail(CodeResultConflict, "重新采购授权已经记录")
|
||||
}
|
||||
now := s.Now()
|
||||
t.RePurchaseAuthorizedAt = &now
|
||||
t.RePurchaseAuthorizedBy = &req.OperatorID
|
||||
t.RePurchaseConsumedAt = nil
|
||||
t.RePurchaseRequestID = &req.RequestID
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Service) ReviewPayment(ctx context.Context, id uint64, req ManualRequest) (models.PurchaseTask, bool, error) {
|
||||
return s.manual(ctx, id, req, func(t *models.PurchaseTask) error {
|
||||
if t.Status != models.PurchaseTaskStatusOrderCreated {
|
||||
return fail(CodeStateConflict, "只有已创建订单才能复核支付状态")
|
||||
}
|
||||
if req.Status != models.PurchasePaymentReviewPaid && req.Status != models.PurchasePaymentReviewUnpaid {
|
||||
return fail(CodeInvalidRequest, "支付复核结果只支持已支付或未支付")
|
||||
}
|
||||
if t.PaymentReviewRequestID != nil {
|
||||
return fail(CodeResultConflict, "支付复核结果已经记录")
|
||||
}
|
||||
now := s.Now()
|
||||
t.PaymentReviewStatus = req.Status
|
||||
t.PaymentReviewedAt = &now
|
||||
t.PaymentReviewedBy = &req.OperatorID
|
||||
t.PaymentReviewRequestID = &req.RequestID
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Service) SelectWriteback(ctx context.Context, id uint64, req ManualRequest) (models.PurchaseTask, bool, error) {
|
||||
var out models.PurchaseTask
|
||||
replayed := false
|
||||
if strings.TrimSpace(req.RequestID) == "" || req.OperatorID == 0 {
|
||||
return out, false, fail(CodeInvalidRequest, "requestId 或操作人无效")
|
||||
}
|
||||
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if e := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&out, id).Error; e != nil {
|
||||
return purchaseNotFound(e)
|
||||
}
|
||||
if out.WritebackSelectRequestID != nil && *out.WritebackSelectRequestID == req.RequestID {
|
||||
replayed = true
|
||||
return nil
|
||||
}
|
||||
if out.Status != models.PurchaseTaskStatusOrderCreated || out.PaymentReviewStatus != models.PurchasePaymentReviewPaid {
|
||||
return fail(CodeStateConflict, "只有已支付订单可以加入回填候选")
|
||||
}
|
||||
if out.SYBProductID != nil {
|
||||
if e := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("syb_product_id = ? AND id <> ? AND writeback_status IN ?", *out.SYBProductID, out.ID, []string{models.PurchaseWritebackStatusPending, models.PurchaseWritebackStatusRunning, models.PurchaseWritebackStatusSucceeded, models.PurchaseWritebackStatusFailed}).Updates(map[string]any{"writeback_status": models.PurchaseWritebackStatusNotSelected, "writeback_select_request_id": nil}).Error; e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
}
|
||||
out.WritebackStatus = models.PurchaseWritebackStatusPending
|
||||
out.WritebackSelectRequestID = &req.RequestID
|
||||
return tx.Save(&out).Error
|
||||
})
|
||||
return out, replayed, err
|
||||
}
|
||||
|
||||
func (s *Service) Cancel(ctx context.Context, id uint64, req ManualRequest) (models.PurchaseTask, bool, error) {
|
||||
return s.manual(ctx, id, req, func(t *models.PurchaseTask) error {
|
||||
if t.Status == models.PurchaseTaskStatusCancelled {
|
||||
return fail(CodeStateConflict, "任务已经取消")
|
||||
}
|
||||
if t.Status == models.PurchaseTaskStatusRunning || t.Status == models.PurchaseTaskStatusOrderSubmitStarted {
|
||||
return fail(CodeStateConflict, "设备正在执行,不能直接取消")
|
||||
}
|
||||
if strings.TrimSpace(req.Reason) == "" {
|
||||
return fail(CodeInvalidRequest, "请填写取消原因")
|
||||
}
|
||||
now := s.Now()
|
||||
if e := t.SetStatus(models.PurchaseTaskStatusCancelled); e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
t.CancelRequestID = &req.RequestID
|
||||
t.CancelledAt = &now
|
||||
t.CancelledBy = &req.OperatorID
|
||||
t.CancelReason = &req.Reason
|
||||
t.StatusVersion++
|
||||
t.StatusChangedAt = now
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Service) ResolveUnknown(ctx context.Context, id uint64, req ManualRequest) (models.PurchaseTask, bool, error) {
|
||||
return s.manual(ctx, id, req, func(t *models.PurchaseTask) error {
|
||||
if t.Status != models.PurchaseTaskStatusOrderResultUnknown {
|
||||
return fail(CodeStateConflict, "任务不是订单结果未知状态")
|
||||
}
|
||||
if req.Status != models.PurchaseTaskStatusOrderCreated && req.Status != models.PurchaseTaskStatusCancelled {
|
||||
return fail(CodeInvalidRequest, "人工处理结果只支持已创建订单或已取消")
|
||||
}
|
||||
now := s.Now()
|
||||
if req.Status == models.PurchaseTaskStatusOrderCreated {
|
||||
if strings.TrimSpace(req.PDDOrderNo) == "" || req.OrderSubmittedAt == nil {
|
||||
return fail(CodeInvalidRequest, "请填写订单号和下单时间")
|
||||
}
|
||||
t.PDDOrderNo = &req.PDDOrderNo
|
||||
t.OrderSubmittedAt = req.OrderSubmittedAt
|
||||
} else {
|
||||
t.CancelledAt = &now
|
||||
t.CancelledBy = &req.OperatorID
|
||||
t.CancelReason = &req.Reason
|
||||
}
|
||||
if e := t.SetStatus(req.Status); e != nil {
|
||||
return internal(e)
|
||||
}
|
||||
t.UnknownResolveRequestID = &req.RequestID
|
||||
t.StatusVersion++
|
||||
t.StatusChangedAt = now
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Service) manual(ctx context.Context, id uint64, req ManualRequest, apply func(*models.PurchaseTask) error) (models.PurchaseTask, bool, error) {
|
||||
var out models.PurchaseTask
|
||||
replayed := false
|
||||
if strings.TrimSpace(req.RequestID) == "" || req.OperatorID == 0 {
|
||||
return out, false, fail(CodeInvalidRequest, "requestId 或操作人无效")
|
||||
}
|
||||
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if e := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&out, id).Error; e != nil {
|
||||
if errors.Is(e, gorm.ErrRecordNotFound) {
|
||||
return fail(CodeTaskNotFound, "采购任务不存在")
|
||||
}
|
||||
return internal(e)
|
||||
}
|
||||
if requestAlreadyApplied(out, req.RequestID) {
|
||||
replayed = true
|
||||
return nil
|
||||
}
|
||||
if e := apply(&out); e != nil {
|
||||
return e
|
||||
}
|
||||
return tx.Save(&out).Error
|
||||
})
|
||||
return out, replayed, err
|
||||
}
|
||||
|
||||
func requestAlreadyApplied(task models.PurchaseTask, requestID string) bool {
|
||||
for _, value := range []*string{task.RePurchaseRequestID, task.PaymentReviewRequestID, task.WritebackSelectRequestID, task.CancelRequestID, task.UnknownResolveRequestID} {
|
||||
if value != nil && *value == requestID {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"os"
|
||||
"strconv"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/common/middleware"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/go-admin-team/go-admin-core/sdk/config"
|
||||
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
|
||||
)
|
||||
|
||||
func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) {
|
||||
h := Handler{}
|
||||
trust, _ := strconv.ParseBool(os.Getenv("GOAUTO_TRUST_FORWARDED_PROTO"))
|
||||
agent := engine.Group("/api/agent/v1/purchase-tasks").Use(device.RequireHTTPS(config.ApplicationConfig.Mode == "prod", trust))
|
||||
agent.GET("/next", h.Next)
|
||||
agent.POST("/:taskId/claim", h.Claim)
|
||||
agent.POST("/:taskId/start", h.Start)
|
||||
agent.POST("/:taskId/order-submit-started", h.OrderSubmitStarted)
|
||||
agent.POST("/:taskId/result", h.Result)
|
||||
admin := engine.Group("/api/admin/v1/purchase-tasks").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole())
|
||||
admin.POST("", h.AdminCreate)
|
||||
admin.POST("/:taskId/spec-decision", h.SpecDecision)
|
||||
admin.POST("/:taskId/authorize-repurchase", h.AuthorizeRePurchase)
|
||||
admin.POST("/:taskId/payment-review", h.ReviewPayment)
|
||||
admin.POST("/:taskId/writeback-candidate", h.SelectWriteback)
|
||||
admin.POST("/:taskId/cancel", h.Cancel)
|
||||
admin.POST("/:taskId/resolve-unknown", h.ResolveUnknown)
|
||||
}
|
||||
@@ -0,0 +1,236 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/app/goauto/purchasecontract"
|
||||
"go-admin/app/goauto/shopeeproduct"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
const DefaultLeaseDuration = 2 * time.Minute
|
||||
|
||||
type Service struct {
|
||||
DB *gorm.DB
|
||||
Now func() time.Time
|
||||
LeaseDuration time.Duration
|
||||
}
|
||||
|
||||
func NewService(db *gorm.DB) *Service {
|
||||
return &Service{DB: db, Now: func() time.Time { return time.Now().UTC() }, LeaseDuration: DefaultLeaseDuration}
|
||||
}
|
||||
|
||||
func (s *Service) Create(ctx context.Context, req CreateRequest) (models.PurchaseTask, bool, error) {
|
||||
if _, err := uuid.Parse(req.RequestID); err != nil {
|
||||
return models.PurchaseTask{}, false, fail(CodeInvalidRequest, "requestId 无效")
|
||||
}
|
||||
if len(req.RuleSnapshot) == 0 {
|
||||
return models.PurchaseTask{}, false, fail(CodeInvalidRequest, "请选择采购规则")
|
||||
}
|
||||
rule, err := purchasecontract.Validate(req.RuleSnapshot, req.ExecutionMode)
|
||||
if err != nil {
|
||||
return models.PurchaseTask{}, false, fail(CodeInvalidRequest, err.Error())
|
||||
}
|
||||
var out models.PurchaseTask
|
||||
replayed := false
|
||||
err = s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Where("create_request_id = ?", req.RequestID).First(&out).Error; err == nil {
|
||||
replayed = true
|
||||
return nil
|
||||
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return internal(err)
|
||||
}
|
||||
var pdd models.PDDProduct
|
||||
var syb models.SYBProduct
|
||||
var shopee models.ShopeeProduct
|
||||
if req.ExecutionMode == models.PurchaseExecutionModeLive {
|
||||
if req.SYBProductID == nil {
|
||||
return fail(CodeInvalidRequest, "正式采购必须选择顺云宝商品")
|
||||
}
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&syb, *req.SYBProductID).Error; err != nil {
|
||||
return notFound(err, "顺云宝商品不存在")
|
||||
}
|
||||
if syb.ShopeeProductID == nil {
|
||||
return fail(CodeInvalidRequest, "该商品尚未关联蝦皮商品")
|
||||
}
|
||||
if err := tx.First(&shopee, *syb.ShopeeProductID).Error; err != nil {
|
||||
return notFound(err, "蝦皮商品不存在")
|
||||
}
|
||||
if shopee.PDDProductID == nil {
|
||||
return fail(CodeInvalidRequest, "该蝦皮商品尚未关联拼多多商品")
|
||||
}
|
||||
if err := tx.First(&pdd, *shopee.PDDProductID).Error; err != nil {
|
||||
return notFound(err, "拼多多商品不存在")
|
||||
}
|
||||
var previous models.PurchaseTask
|
||||
previousErr := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("syb_product_id = ?", syb.ID).Order("id DESC").First(&previous).Error
|
||||
if previousErr == nil {
|
||||
switch previous.Status {
|
||||
case models.PurchaseTaskStatusFailed, models.PurchaseTaskStatusCancelled, models.PurchaseTaskStatusRehearsalCompleted:
|
||||
case models.PurchaseTaskStatusOrderCreated:
|
||||
if previous.RePurchaseAuthorizedAt == nil || previous.RePurchaseConsumedAt != nil {
|
||||
return fail(CodeRePurchaseRequired, "该商品已经采购成功,需要先授权重新采购")
|
||||
}
|
||||
now := s.Now()
|
||||
if err := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ? AND re_purchase_consumed_at IS NULL", previous.ID).Updates(map[string]any{"re_purchase_consumed_at": now, "active_slot": nil}).Error; err != nil {
|
||||
return internal(err)
|
||||
}
|
||||
default:
|
||||
return fail(CodeStateConflict, "该顺云宝商品已有未结束的采购任务")
|
||||
}
|
||||
} else if !errors.Is(previousErr, gorm.ErrRecordNotFound) {
|
||||
return internal(previousErr)
|
||||
}
|
||||
} else if req.ExecutionMode == models.PurchaseExecutionModeRehearsal {
|
||||
if req.PDDProductID == nil {
|
||||
return fail(CodeInvalidRequest, "演练必须选择拼多多商品")
|
||||
}
|
||||
if err := tx.First(&pdd, *req.PDDProductID).Error; err != nil {
|
||||
return notFound(err, "拼多多商品不存在")
|
||||
}
|
||||
} else {
|
||||
return fail(CodeInvalidRequest, "executionMode 只支持 rehearsal 或 live")
|
||||
}
|
||||
if req.DeviceID != nil {
|
||||
var d models.AgentDevice
|
||||
if err := tx.First(&d, *req.DeviceID).Error; err != nil {
|
||||
return notFound(err, "设备不存在")
|
||||
}
|
||||
if err := ensureCapabilities(d, purchasecontract.RequiredCapabilities(rule)); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if req.PDDAccountID != nil {
|
||||
var a models.PDDAccount
|
||||
if err := tx.First(&a, *req.PDDAccountID).Error; err != nil {
|
||||
return notFound(err, "拼多多账号不存在")
|
||||
}
|
||||
if a.Status != "active" {
|
||||
return fail(CodeInvalidRequest, "拼多多账号已停用")
|
||||
}
|
||||
}
|
||||
quantity, ref, minPrice, maxPrice, currency := req.Quantity, req.ReferenceUnitPriceCent, req.MinUnitPriceCent, req.MaxUnitPriceCent, strings.ToUpper(strings.TrimSpace(req.Currency))
|
||||
targetColor, targetSize, mappedColor, mappedSize, specSource := strings.TrimSpace(req.TargetColor), strings.TrimSpace(req.TargetSize), strings.TrimSpace(req.MappedColor), strings.TrimSpace(req.MappedSize), "unresolved"
|
||||
var sybID, shopeeID *uint64
|
||||
if req.ExecutionMode == models.PurchaseExecutionModeLive {
|
||||
sybID, shopeeID = &syb.ID, &shopee.ID
|
||||
quantity, ref = syb.Quantity, syb.UnitPriceCent
|
||||
if targetColor == "" {
|
||||
targetColor = syb.TargetColor
|
||||
}
|
||||
if targetSize == "" {
|
||||
targetSize = syb.TargetSize
|
||||
}
|
||||
if currency == "" {
|
||||
currency = shopee.Currency
|
||||
}
|
||||
mappedColor, mappedSize, specSource = confirmedMappings(shopee.SpecsJSON, targetColor, targetSize)
|
||||
} else if (targetColor == "" || mappedColor != "") && (targetSize == "" || mappedSize != "") {
|
||||
specSource = "manual_mapping"
|
||||
}
|
||||
if pdd.Status == "pending" || strings.TrimSpace(pdd.SpecsJSON) == "" || strings.TrimSpace(pdd.SpecsJSON) == "[]" {
|
||||
mappedColor, mappedSize, specSource = "", "", "unresolved"
|
||||
}
|
||||
if specSource == "unresolved" && !containsString(rule.RequiredCapabilities, purchasecontract.CapabilitySpecProbeV1) {
|
||||
return fail(CodeMappingRequired, "规格映射不完整,所选规则不支持规格探测")
|
||||
}
|
||||
if quantity < 1 || minPrice < 0 || maxPrice < minPrice || currency == "" {
|
||||
return fail(CodeInvalidRequest, "数量、价格区间或币种无效")
|
||||
}
|
||||
required, _ := json.Marshal(purchasecontract.RequiredCapabilities(rule))
|
||||
out = models.PurchaseTask{SYBProductID: sybID, ShopeeProductID: shopeeID, PDDProductID: pdd.ID, DeviceID: req.DeviceID, PDDAccountID: req.PDDAccountID,
|
||||
ExecutionMode: req.ExecutionMode, Status: models.PurchaseTaskStatusPending, ShopeeItemIDSnapshot: shopee.ShopeeItemID, ShopeeTitleSnapshot: shopee.Title, ShopeeShopNameSnapshot: shopee.ShopName,
|
||||
PDDURLSnapshot: pdd.URL, PDDGoodsIDSnapshot: pdd.GoodsID, PDDTitleSnapshot: pdd.Title, TargetColorSnapshot: targetColor, TargetSizeSnapshot: targetSize, MappedColorSnapshot: mappedColor, MappedSizeSnapshot: mappedSize, SpecSource: specSource,
|
||||
Quantity: quantity, ReferenceUnitPriceCent: ref, MinUnitPriceCent: minPrice, MaxUnitPriceCent: maxPrice, Currency: currency, RuleType: rule.RuleType, RuleSchemaVersion: rule.SchemaVersion, RequiredCapabilitiesJSON: string(required), RuleSnapshot: string(req.RuleSnapshot), CreateRequestID: req.RequestID,
|
||||
PaymentReviewStatus: models.PurchasePaymentReviewPending, LogisticsStatus: models.PurchaseLogisticsStatusPending, WritebackStatus: models.PurchaseWritebackStatusNotSelected}
|
||||
if err := tx.Create(&out).Error; err != nil {
|
||||
if isDuplicate(err) {
|
||||
return fail(CodeRePurchaseRequired, "该顺云宝商品已有采购任务,不能重复创建")
|
||||
}
|
||||
return internal(err)
|
||||
}
|
||||
out.AddressSuffix = purchasecontract.AddressSuffix(out.ID)
|
||||
if err := tx.Model(&out).Update("address_suffix", out.AddressSuffix).Error; err != nil {
|
||||
return internal(err)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return out, replayed, err
|
||||
}
|
||||
|
||||
func confirmedMappings(raw, color, size string) (string, string, string) {
|
||||
var specs []shopeeproduct.SpecDimension
|
||||
if json.Unmarshal([]byte(raw), &specs) != nil {
|
||||
return "", "", "unresolved"
|
||||
}
|
||||
mappedColor, mappedSize, source := "", "", "manual_mapping"
|
||||
for _, d := range specs {
|
||||
for _, v := range d.Values {
|
||||
if v.Mapping == nil || v.Mapping.Status != shopeeproduct.MappingStatusConfirmed {
|
||||
continue
|
||||
}
|
||||
if d.Role == shopeeproduct.RoleColor && v.Name == color {
|
||||
mappedColor = v.Mapping.PDDValue
|
||||
source = mapSource(v.Mapping.Source)
|
||||
}
|
||||
if d.Role == shopeeproduct.RoleSize && v.Name == size {
|
||||
mappedSize = v.Mapping.PDDValue
|
||||
source = mapSource(v.Mapping.Source)
|
||||
}
|
||||
}
|
||||
}
|
||||
if (color != "" && mappedColor == "") || (size != "" && mappedSize == "") {
|
||||
return mappedColor, mappedSize, "unresolved"
|
||||
}
|
||||
return mappedColor, mappedSize, source
|
||||
}
|
||||
func mapSource(v string) string {
|
||||
if v == shopeeproduct.MappingSourceExactMatch {
|
||||
return "exact_match"
|
||||
}
|
||||
if v == shopeeproduct.MappingSourceAIMatch {
|
||||
return "ai_match"
|
||||
}
|
||||
return "manual_mapping"
|
||||
}
|
||||
func notFound(err error, msg string) error {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return fail(CodeInvalidRequest, msg)
|
||||
}
|
||||
return internal(err)
|
||||
}
|
||||
func isDuplicate(err error) bool {
|
||||
v := strings.ToLower(err.Error())
|
||||
return strings.Contains(v, "duplicate") || strings.Contains(v, "unique constraint")
|
||||
}
|
||||
func ensureCapabilities(d models.AgentDevice, required []string) error {
|
||||
ok, err := device.Supports(d, required)
|
||||
if err != nil {
|
||||
return internal(err)
|
||||
}
|
||||
if !ok {
|
||||
return fail(CodeCapabilityMismatch, "设备能力不符合采购规则")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func hashBytes(raw []byte) string { h := sha256.Sum256(raw); return hex.EncodeToString(h[:]) }
|
||||
|
||||
func containsString(values []string, target string) bool {
|
||||
for _, value := range values {
|
||||
if value == target {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -0,0 +1,325 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/app/goauto/migrations"
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/app/goauto/purchasecontract"
|
||||
"go-admin/app/goauto/shopeeproduct"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
type fixture struct {
|
||||
syb models.SYBProduct
|
||||
shopee models.ShopeeProduct
|
||||
pdd models.PDDProduct
|
||||
device models.AgentDevice
|
||||
token string
|
||||
}
|
||||
|
||||
func testDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
dsn := fmt.Sprintf("file:%s?mode=memory&cache=shared&_foreign_keys=on", strings.ReplaceAll(t.Name(), "/", "_"))
|
||||
db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err = migrations.Migrate(db); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return db
|
||||
}
|
||||
func testService(db *gorm.DB) *Service {
|
||||
s := NewService(db)
|
||||
s.Now = func() time.Time { return time.Date(2026, 8, 20, 12, 0, 0, 0, time.UTC) }
|
||||
return s
|
||||
}
|
||||
func seed(t *testing.T, db *gorm.DB, caps []string, mapped bool) fixture {
|
||||
t.Helper()
|
||||
p := models.PDDProduct{GoodsID: "719834019024", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=719834019024", Title: "测试商品", Status: "active", SpecsJSON: `[{"name":"颜色","role":"color","values":[{"name":"黑色"}]}]`}
|
||||
if err := db.Create(&p).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
mapping := (*shopeeproduct.Mapping)(nil)
|
||||
if mapped {
|
||||
mapping = &shopeeproduct.Mapping{PDDValue: "黑色", Source: shopeeproduct.MappingSourceManual, Status: shopeeproduct.MappingStatusConfirmed}
|
||||
}
|
||||
specs := []shopeeproduct.SpecDimension{{Name: "颜色", Role: shopeeproduct.RoleColor, Values: []shopeeproduct.SpecValue{{Name: "黑色", Source: shopeeproduct.ValueSourceImport, Mapping: mapping}}}, {Name: "尺码", Role: shopeeproduct.RoleSize, Values: []shopeeproduct.SpecValue{{Name: "XL", Source: shopeeproduct.ValueSourceImport, Mapping: &shopeeproduct.Mapping{PDDValue: "XL", Source: shopeeproduct.MappingSourceManual, Status: shopeeproduct.MappingStatusConfirmed}}}}}
|
||||
raw, _ := json.Marshal(specs)
|
||||
sp := models.ShopeeProduct{ShopeeItemID: "26154802794", Title: "蝦皮商品", ShopName: "测试店", PDDProductID: &p.ID, SpecsJSON: string(raw), Currency: "CNY"}
|
||||
if err := db.Create(&sp).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
syb := models.SYBProduct{OrderCode: "SYB-1", DetailID: 1, StockID: 2, ShopeeItemID: sp.ShopeeItemID, ShopeeProductID: &sp.ID, ProductTitle: sp.Title, TargetColor: "黑色", TargetSize: "XL", Quantity: 2, UnitPriceCent: 2000, ImageURL: "", ParseStatus: models.SYBParseStatusSuccess, RawJSON: `{}`}
|
||||
if err := db.Create(&syb).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
token := uuid.NewString()
|
||||
ds := device.NewService(db)
|
||||
ds.GenerateToken = func() (string, error) { return token, nil }
|
||||
response, err := ds.Register(context.Background(), device.RegisterRequest{RequestID: uuid.NewString(), InstallID: uuid.NewString(), Name: "Samsung", Manufacturer: "Samsung", Model: "S24", AndroidVersion: "15", AgentVersion: "1", PDDVersion: "7", Capabilities: caps}, "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var d models.AgentDevice
|
||||
if err = db.First(&d, response.DeviceID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return fixture{syb: syb, shopee: sp, pdd: p, device: d, token: token}
|
||||
}
|
||||
func liveCaps() []string {
|
||||
return []string{purchasecontract.CapabilityPurchaseLiveV1, purchasecontract.CapabilityAddressUpdateV1, purchasecontract.CapabilityOrderCreateV1, purchasecontract.CapabilitySpecProbeV1}
|
||||
}
|
||||
func liveRule(probe bool) []byte {
|
||||
cap := `"purchase.live.v1","purchase.address-update.v1","purchase.order-create.v1"`
|
||||
actions := `{"type":"openProduct"},{"type":"updateShippingAddress"},{"type":"createOrder"},{"type":"readOrderResult"}`
|
||||
if probe {
|
||||
cap += `,"purchase.spec-probe.v1"`
|
||||
actions += `,{"type":"probeSpecs"}`
|
||||
}
|
||||
return []byte(`{"schemaVersion":1,"ruleType":"pddPurchase","requiredCapabilities":[` + cap + `],"actions":[` + actions + `]}`)
|
||||
}
|
||||
func createLive(t *testing.T, s *Service, f fixture) (models.PurchaseTask, error) {
|
||||
t.Helper()
|
||||
r, _, err := s.Create(context.Background(), CreateRequest{RequestID: uuid.NewString(), ExecutionMode: models.PurchaseExecutionModeLive, SYBProductID: &f.syb.ID, DeviceID: &f.device.ID, MinUnitPriceCent: 400, MaxUnitPriceCent: 3000, RuleSnapshot: liveRule(true)})
|
||||
return r, err
|
||||
}
|
||||
func code(err error) string {
|
||||
var e *ServiceError
|
||||
if errors.As(err, &e) {
|
||||
return e.Code
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func TestRehearsalCannotContainOrderActions(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, []string{purchasecontract.CapabilityPurchaseRehearsalV1}, true)
|
||||
_, _, err := testService(db).Create(context.Background(), CreateRequest{RequestID: uuid.NewString(), ExecutionMode: models.PurchaseExecutionModeRehearsal, PDDProductID: &f.pdd.ID, Quantity: 1, Currency: "CNY", MaxUnitPriceCent: 1000, RuleSnapshot: []byte(`{"schemaVersion":1,"ruleType":"pddPurchase","requiredCapabilities":["purchase.rehearsal.v1","purchase.order-create.v1"],"actions":[{"type":"createOrder"}]}`)})
|
||||
if code(err) != CodeInvalidRequest {
|
||||
t.Fatalf("expected rehearsal rejection, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCreateAndLifecycleValidateCapabilitiesAndIdempotentResult(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
task, err := createLive(t, s, f)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if task.AddressSuffix != "_cg1" {
|
||||
t.Fatalf("unexpected suffix %s", task.AddressSuffix)
|
||||
}
|
||||
if _, err = s.Claim(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
start, err := s.Start(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if start.Phase != models.PurchaseAttemptPhasePurchase {
|
||||
t.Fatalf("unexpected phase %s", start.Phase)
|
||||
}
|
||||
if _, err = s.MarkOrderSubmitStarted(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
submitted := s.Now()
|
||||
req := ResultRequest{RequestID: uuid.NewString(), TaskAttemptID: start.TaskAttemptID, ResultType: "order_created", PDDOrderNo: "PDD-1", OrderSubmittedAt: &submitted}
|
||||
first, err := s.SubmitResult(context.Background(), task.ID, req, f.token)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
replay, err := s.SubmitResult(context.Background(), task.ID, req, f.token)
|
||||
if err != nil || !replay.Replayed {
|
||||
t.Fatalf("result replay failed: %+v %v", replay, err)
|
||||
}
|
||||
var count int64
|
||||
db.Model(&models.PurchaseTaskAttempt{}).Where("task_id = ?", task.ID).Count(&count)
|
||||
if count != 1 || first.Status != models.PurchaseTaskStatusOrderCreated {
|
||||
t.Fatalf("duplicate attempt or wrong status: %d %+v", count, first)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCapabilityMismatchAtCreateAndClaim(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, []string{purchasecontract.CapabilityPurchaseLiveV1}, true)
|
||||
_, err := createLive(t, testService(db), f)
|
||||
if code(err) != CodeCapabilityMismatch {
|
||||
t.Fatalf("create should reject capability mismatch: %v", err)
|
||||
}
|
||||
f.device.CapabilitiesJSON = `["purchase.live.v1","purchase.address-update.v1","purchase.order-create.v1","purchase.spec-probe.v1"]`
|
||||
db.Model(&models.AgentDevice{}).Where("id = ?", f.device.ID).Update("capabilities_json", f.device.CapabilitiesJSON)
|
||||
task, err := createLive(t, testService(db), f)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
db.Model(&models.AgentDevice{}).Where("id = ?", f.device.ID).Update("capabilities_json", `["purchase.live.v1"]`)
|
||||
_, err = testService(db).Claim(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
if code(err) != CodeCapabilityMismatch {
|
||||
t.Fatalf("claim should reject changed capabilities: %v", err)
|
||||
}
|
||||
db.Model(&models.AgentDevice{}).Where("id = ?", f.device.ID).Update("capabilities_json", f.device.CapabilitiesJSON)
|
||||
if _, err = testService(db).Claim(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
db.Model(&models.AgentDevice{}).Where("id = ?", f.device.ID).Update("capabilities_json", `["purchase.live.v1"]`)
|
||||
_, err = testService(db).Start(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
if code(err) != CodeCapabilityMismatch {
|
||||
t.Fatalf("start should reject changed capabilities: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSlowPathUsesTwoAttemptsAndFreezesDecision(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), false)
|
||||
s := testService(db)
|
||||
task, err := createLive(t, s, f)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err = s.Claim(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
first, err := s.Start(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
if err != nil || first.Phase != models.PurchaseAttemptPhaseSpecProbe {
|
||||
t.Fatalf("probe start: %+v %v", first, err)
|
||||
}
|
||||
probe := ResultRequest{RequestID: uuid.NewString(), TaskAttemptID: first.TaskAttemptID, ResultType: "spec_probe_completed", ProbedSpecs: []byte(`[{"color":"黑色","size":"XL"}]`)}
|
||||
if _, err = s.SubmitResult(context.Background(), task.ID, probe, f.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
decision := SpecDecisionRequest{RequestID: uuid.NewString(), TaskAttemptID: first.TaskAttemptID, MappedColor: "黑色", MappedSize: "XL", Source: "ai_match", Decision: []byte(`{"reason":"same label"}`), OperatorID: 1}
|
||||
if _, _, err = s.ApplySpecDecision(context.Background(), task.ID, decision); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
decision.MappedColor = "白色"
|
||||
if _, _, err = s.ApplySpecDecision(context.Background(), task.ID, decision); code(err) != CodeResultConflict {
|
||||
t.Fatalf("frozen decision changed: %v", err)
|
||||
}
|
||||
if _, err = s.Claim(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
second, err := s.Start(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
if err != nil || second.Phase != models.PurchaseAttemptPhasePurchase || second.AttemptNumber != 2 {
|
||||
t.Fatalf("second attempt: %+v %v", second, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderUnknownIsNotAutomaticallyRedispatched(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
task, _ := createLive(t, s, f)
|
||||
s.Claim(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
started, _ := s.Start(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
s.MarkOrderSubmitStarted(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, f.token)
|
||||
if _, err := s.SubmitResult(context.Background(), task.ID, ResultRequest{RequestID: uuid.NewString(), TaskAttemptID: started.TaskAttemptID, ResultType: "order_result_unknown"}, f.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
next, err := s.Next(context.Background(), f.token)
|
||||
if err != nil || next != nil {
|
||||
t.Fatalf("unknown task redispatched: %+v %v", next, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestClaimLocksKnownPDDAccountAcrossDevices(t *testing.T) {
|
||||
db := testDB(t)
|
||||
first := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
account := models.PDDAccount{Name: "account-a", Status: "active"}
|
||||
if err := db.Create(&account).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
secondSYB := first.syb
|
||||
secondSYB.ID = 0
|
||||
secondSYB.OrderCode = "SYB-2"
|
||||
secondSYB.DetailID = 2
|
||||
if err := db.Create(&secondSYB).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
token := uuid.NewString()
|
||||
ds := device.NewService(db)
|
||||
ds.GenerateToken = func() (string, error) { return token, nil }
|
||||
registered, err := ds.Register(context.Background(), device.RegisterRequest{RequestID: uuid.NewString(), InstallID: uuid.NewString(), Name: "OnePlus", Manufacturer: "OnePlus", Model: "12", AndroidVersion: "15", AgentVersion: "1", PDDVersion: "7", Capabilities: liveCaps()}, "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var secondDevice models.AgentDevice
|
||||
db.First(&secondDevice, registered.DeviceID)
|
||||
create := func(sybID, deviceID uint64) (models.PurchaseTask, error) {
|
||||
record, _, err := s.Create(context.Background(), CreateRequest{RequestID: uuid.NewString(), ExecutionMode: models.PurchaseExecutionModeLive, SYBProductID: &sybID, DeviceID: &deviceID, PDDAccountID: &account.ID, MinUnitPriceCent: 400, MaxUnitPriceCent: 3000, RuleSnapshot: liveRule(true)})
|
||||
return record, err
|
||||
}
|
||||
firstTask, err := create(first.syb.ID, first.device.ID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
secondTask, err := create(secondSYB.ID, secondDevice.ID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err = s.Claim(context.Background(), firstTask.ID, ActionRequest{RequestID: uuid.NewString()}, first.token); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err = s.Claim(context.Background(), secondTask.ID, ActionRequest{RequestID: uuid.NewString()}, token); code(err) != CodeDeviceBusy {
|
||||
t.Fatalf("second account claim should fail: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSuccessfulTaskNeedsOneTimeRePurchaseAuthorization(t *testing.T) {
|
||||
db := testDB(t)
|
||||
f := seed(t, db, liveCaps(), true)
|
||||
s := testService(db)
|
||||
task, err := createLive(t, s, f)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
now := s.Now()
|
||||
task.SetStatus(models.PurchaseTaskStatusOrderCreated)
|
||||
task.PDDOrderNo = ptrString("PDD-1")
|
||||
task.OrderSubmittedAt = &now
|
||||
if err = db.Save(&task).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err = createLive(t, s, f); code(err) != CodeRePurchaseRequired {
|
||||
t.Fatalf("expected authorization requirement: %v", err)
|
||||
}
|
||||
auth := ManualRequest{RequestID: uuid.NewString(), OperatorID: 1}
|
||||
if _, _, err = s.AuthorizeRePurchase(context.Background(), task.ID, auth); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
second, err := createLive(t, s, f)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var old models.PurchaseTask
|
||||
db.First(&old, task.ID)
|
||||
if old.RePurchaseConsumedAt == nil || second.ID == task.ID {
|
||||
t.Fatalf("authorization not consumed: %+v", old)
|
||||
}
|
||||
if _, _, err = s.ReviewPayment(context.Background(), old.ID, ManualRequest{RequestID: uuid.NewString(), OperatorID: 1, Status: models.PurchasePaymentReviewPaid}); err != nil {
|
||||
t.Fatalf("old order must remain editable after authorization consumption: %v", err)
|
||||
}
|
||||
second.SetStatus(models.PurchaseTaskStatusOrderCreated)
|
||||
second.PDDOrderNo = ptrString("PDD-2")
|
||||
second.OrderSubmittedAt = &now
|
||||
db.Save(&second)
|
||||
if _, err = createLive(t, s, f); code(err) != CodeRePurchaseRequired {
|
||||
t.Fatalf("authorization reused: %v", err)
|
||||
}
|
||||
}
|
||||
func ptrString(v string) *string { return &v }
|
||||
@@ -0,0 +1,121 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
CodeInvalidRequest = "PURCHASE_INVALID_REQUEST"
|
||||
CodeTaskNotFound = "PURCHASE_TASK_NOT_FOUND"
|
||||
CodeStateConflict = "PURCHASE_STATE_CONFLICT"
|
||||
CodeCapabilityMismatch = "PURCHASE_CAPABILITY_MISMATCH"
|
||||
CodeDeviceBusy = "DEVICE_BUSY"
|
||||
CodeTaskClaimed = "PURCHASE_TASK_ALREADY_CLAIMED"
|
||||
CodeLeaseExpired = "PURCHASE_LEASE_EXPIRED"
|
||||
CodeMappingRequired = "PURCHASE_SPEC_MAPPING_REQUIRED"
|
||||
CodeResultConflict = "PURCHASE_RESULT_CONFLICT"
|
||||
CodeRePurchaseRequired = "REPURCHASE_AUTHORIZATION_REQUIRED"
|
||||
CodeInternal = "INTERNAL_ERROR"
|
||||
)
|
||||
|
||||
type ServiceError struct {
|
||||
Code string
|
||||
Message string
|
||||
Retryable bool
|
||||
Cause error
|
||||
}
|
||||
|
||||
func (e *ServiceError) Error() string {
|
||||
if e.Cause == nil {
|
||||
return e.Message
|
||||
}
|
||||
return fmt.Sprintf("%s: %v", e.Message, e.Cause)
|
||||
}
|
||||
func (e *ServiceError) Unwrap() error { return e.Cause }
|
||||
func fail(code, message string) error { return &ServiceError{Code: code, Message: message} }
|
||||
func internal(err error) error {
|
||||
return &ServiceError{Code: CodeInternal, Message: "服务端处理失败", Retryable: true, Cause: err}
|
||||
}
|
||||
|
||||
type CreateRequest struct {
|
||||
RequestID string `json:"requestId"`
|
||||
ExecutionMode string `json:"executionMode"`
|
||||
SYBProductID *uint64 `json:"sybProductId,omitempty"`
|
||||
PDDProductID *uint64 `json:"pddProductId,omitempty"`
|
||||
DeviceID *uint64 `json:"deviceId,omitempty"`
|
||||
PDDAccountID *uint64 `json:"pddAccountId,omitempty"`
|
||||
TargetColor string `json:"targetColor,omitempty"`
|
||||
TargetSize string `json:"targetSize,omitempty"`
|
||||
MappedColor string `json:"mappedColor,omitempty"`
|
||||
MappedSize string `json:"mappedSize,omitempty"`
|
||||
Quantity int64 `json:"quantity,omitempty"`
|
||||
ReferenceUnitPriceCent int64 `json:"referenceUnitPriceCent,omitempty"`
|
||||
MinUnitPriceCent int64 `json:"minUnitPriceCent,omitempty"`
|
||||
MaxUnitPriceCent int64 `json:"maxUnitPriceCent,omitempty"`
|
||||
Currency string `json:"currency,omitempty"`
|
||||
RuleSnapshot json.RawMessage `json:"ruleSnapshot"`
|
||||
}
|
||||
|
||||
type ActionRequest struct {
|
||||
RequestID string `json:"requestId"`
|
||||
}
|
||||
|
||||
type TaskPayload struct {
|
||||
TaskID uint64 `json:"taskId"`
|
||||
TaskAttemptID string `json:"taskAttemptId,omitempty"`
|
||||
AttemptNumber int `json:"attemptNumber,omitempty"`
|
||||
Phase string `json:"phase,omitempty"`
|
||||
ExecutionMode string `json:"executionMode"`
|
||||
Status string `json:"status"`
|
||||
DeviceID *uint64 `json:"deviceId,omitempty"`
|
||||
PDDProductID uint64 `json:"pddProductId"`
|
||||
PDDURL string `json:"pddUrl"`
|
||||
PDDGoodsID string `json:"pddGoodsId"`
|
||||
TargetColor string `json:"targetColor"`
|
||||
TargetSize string `json:"targetSize"`
|
||||
MappedColor string `json:"mappedColor"`
|
||||
MappedSize string `json:"mappedSize"`
|
||||
Quantity int64 `json:"quantity"`
|
||||
MinUnitPriceCent int64 `json:"minUnitPriceCent"`
|
||||
MaxUnitPriceCent int64 `json:"maxUnitPriceCent"`
|
||||
Currency string `json:"currency"`
|
||||
AddressSuffix string `json:"addressSuffix"`
|
||||
RuleSnapshot json.RawMessage `json:"ruleSnapshot"`
|
||||
RuleSnapshotHash string `json:"ruleSnapshotHash,omitempty"`
|
||||
LeaseExpiresAt *time.Time `json:"leaseExpiresAt,omitempty"`
|
||||
LeaseVersion uint64 `json:"leaseVersion"`
|
||||
Replayed bool `json:"replayed,omitempty"`
|
||||
}
|
||||
|
||||
type ResultRequest struct {
|
||||
RequestID string `json:"requestId"`
|
||||
TaskAttemptID string `json:"taskAttemptId"`
|
||||
ResultType string `json:"resultType"`
|
||||
PDDOrderNo string `json:"pddOrderNo,omitempty"`
|
||||
OrderSubmittedAt *time.Time `json:"orderSubmittedAt,omitempty"`
|
||||
ErrorCode string `json:"errorCode,omitempty"`
|
||||
ErrorMessage string `json:"errorMessage,omitempty"`
|
||||
ProbedSpecs json.RawMessage `json:"probedSpecs,omitempty"`
|
||||
}
|
||||
|
||||
type SpecDecisionRequest struct {
|
||||
RequestID string `json:"requestId"`
|
||||
TaskAttemptID string `json:"taskAttemptId"`
|
||||
MappedColor string `json:"mappedColor"`
|
||||
MappedSize string `json:"mappedSize"`
|
||||
Source string `json:"source"`
|
||||
Decision json.RawMessage `json:"decision"`
|
||||
NoMatch bool `json:"noMatch,omitempty"`
|
||||
OperatorID uint64 `json:"-"`
|
||||
}
|
||||
|
||||
type ManualRequest struct {
|
||||
RequestID string `json:"requestId"`
|
||||
OperatorID uint64 `json:"operatorId"`
|
||||
Reason string `json:"reason,omitempty"`
|
||||
Status string `json:"status,omitempty"`
|
||||
PDDOrderNo string `json:"pddOrderNo,omitempty"`
|
||||
OrderSubmittedAt *time.Time `json:"orderSubmittedAt,omitempty"`
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
package version_local
|
||||
|
||||
import (
|
||||
"runtime"
|
||||
|
||||
goautomigrations "go-admin/app/goauto/migrations"
|
||||
"go-admin/cmd/migrate/migration"
|
||||
common "go-admin/common/models"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// This additive migration extends the #33 contract with the idempotency and
|
||||
// audit columns required by the #34 purchase task state machine. It preserves
|
||||
// every existing task, attempt and product row.
|
||||
func init() {
|
||||
_, fileName, _, _ := runtime.Caller(0)
|
||||
migration.Migrate.SetVersion(migration.GetFilename(fileName), migratePurchaseStateMachine)
|
||||
}
|
||||
|
||||
func migratePurchaseStateMachine(db *gorm.DB, version string) error {
|
||||
return db.Transaction(func(tx *gorm.DB) error {
|
||||
if err := goautomigrations.Migrate(tx); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(&common.Migration{Version: version}).Error
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user