From 65ae41fd6fdd9a0ee31b3c7a3fab1d0fec58db6e Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Thu, 20 Aug 2026 17:22:52 +0800 Subject: [PATCH] feat(#34): implement purchase task state machine --- docs/00-project-profile.md | 4 +- docs/02-architecture-and-code-map.md | 3 +- docs/03-business-rules-and-glossary.md | 7 +- docs/08-agent-api-contract.md | 23 +- docs/09-delivery-issues.md | 6 +- server/app/admin/router/init_router.go | 2 + .../migrations/purchase_contract_test.go | 5 +- server/app/goauto/models/purchase.go | 133 +++-- server/app/goauto/purchase/handler.go | 250 +++++++++ server/app/goauto/purchase/lifecycle.go | 488 ++++++++++++++++++ server/app/goauto/purchase/manual.go | 166 ++++++ server/app/goauto/purchase/router.go | 32 ++ server/app/goauto/purchase/service.go | 236 +++++++++ server/app/goauto/purchase/service_test.go | 325 ++++++++++++ server/app/goauto/purchase/types.go | 121 +++++ .../1786701200000_purchase_state_machine.go | 28 + 16 files changed, 1758 insertions(+), 71 deletions(-) create mode 100644 server/app/goauto/purchase/handler.go create mode 100644 server/app/goauto/purchase/lifecycle.go create mode 100644 server/app/goauto/purchase/manual.go create mode 100644 server/app/goauto/purchase/router.go create mode 100644 server/app/goauto/purchase/service.go create mode 100644 server/app/goauto/purchase/service_test.go create mode 100644 server/app/goauto/purchase/types.go create mode 100644 server/cmd/migrate/migration/version-local/1786701200000_purchase_state_machine.go diff --git a/docs/00-project-profile.md b/docs/00-project-profile.md index 69d5788..cddef3e 100644 --- a/docs/00-project-profile.md +++ b/docs/00-project-profile.md @@ -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 @@ -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 订单。 diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index 6abd25f..3733123 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -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 @@ -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/` | diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 5c3bcf5..2632f58 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -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 @@ -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 明细后来选择的订单覆盖旧候选,但不删除旧订单事实。 ## 自动化边界 diff --git a/docs/08-agent-api-contract.md b/docs/08-agent-api-contract.md index 6dc788d..a47d80a 100644 --- a/docs/08-agent-api-contract.md +++ b/docs/08-agent-api-contract.md @@ -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 @@ -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 凭据或完整收货地址。 diff --git a/docs/09-delivery-issues.md b/docs/09-delivery-issues.md index 2eb345d..c8cd087 100644 --- a/docs/09-delivery-issues.md +++ b/docs/09-delivery-issues.md @@ -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 @@ -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) | 服务端物流调度与货运宝自动回填 | 采购任务与有效订单能力 | diff --git a/server/app/admin/router/init_router.go b/server/app/admin/router/init_router.go index 7111b6f..5f92e5a 100644 --- a/server/app/admin/router/init_router.go +++ b/server/app/admin/router/init_router.go @@ -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) diff --git a/server/app/goauto/migrations/purchase_contract_test.go b/server/app/goauto/migrations/purchase_contract_test.go index 50a11ca..e7a9a1a 100644 --- a/server/app/goauto/migrations/purchase_contract_test.go +++ b/server/app/goauto/migrations/purchase_contract_test.go @@ -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) } diff --git a/server/app/goauto/models/purchase.go b/server/app/goauto/models/purchase.go index 5efb882..5547b68 100644 --- a/server/app/goauto/models/purchase.go +++ b/server/app/goauto/models/purchase.go @@ -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" } diff --git a/server/app/goauto/purchase/handler.go b/server/app/goauto/purchase/handler.go new file mode 100644 index 0000000..535ff75 --- /dev/null +++ b/server/app/goauto/purchase/handler.go @@ -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 +} diff --git a/server/app/goauto/purchase/lifecycle.go b/server/app/goauto/purchase/lifecycle.go new file mode 100644 index 0000000..85710a7 --- /dev/null +++ b/server/app/goauto/purchase/lifecycle.go @@ -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 +} diff --git a/server/app/goauto/purchase/manual.go b/server/app/goauto/purchase/manual.go new file mode 100644 index 0000000..fd93083 --- /dev/null +++ b/server/app/goauto/purchase/manual.go @@ -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 +} diff --git a/server/app/goauto/purchase/router.go b/server/app/goauto/purchase/router.go new file mode 100644 index 0000000..b26cf28 --- /dev/null +++ b/server/app/goauto/purchase/router.go @@ -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) +} diff --git a/server/app/goauto/purchase/service.go b/server/app/goauto/purchase/service.go new file mode 100644 index 0000000..69834b8 --- /dev/null +++ b/server/app/goauto/purchase/service.go @@ -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 +} diff --git a/server/app/goauto/purchase/service_test.go b/server/app/goauto/purchase/service_test.go new file mode 100644 index 0000000..3451459 --- /dev/null +++ b/server/app/goauto/purchase/service_test.go @@ -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 } diff --git a/server/app/goauto/purchase/types.go b/server/app/goauto/purchase/types.go new file mode 100644 index 0000000..394445c --- /dev/null +++ b/server/app/goauto/purchase/types.go @@ -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"` +} diff --git a/server/cmd/migrate/migration/version-local/1786701200000_purchase_state_machine.go b/server/cmd/migrate/migration/version-local/1786701200000_purchase_state_machine.go new file mode 100644 index 0000000..25240fc --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1786701200000_purchase_state_machine.go @@ -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 + }) +}