Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0ab5a96ca4 | ||
|
|
5e0a9d108c | ||
|
|
08b7095cf1 | ||
|
|
01510a85dc | ||
|
|
4261a542ca | ||
|
|
beec630187 | ||
|
|
2491a857f7 | ||
|
|
0661b2205f | ||
|
|
bb1a410e8a |
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Project-Profile
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Project-Profile.-
|
||||
wiki_revision: 3b78360779ae520f1ff9e51be3118a04cd549f51
|
||||
synchronized_at: 2026-09-07T09:27:40Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:14:21Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 项目档案
|
||||
|
||||
+2
-2
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Development-Workflow
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Development-Workflow.-
|
||||
wiki_revision: 62ddbe4469740c02ce4a6ca2fd1966a89a79322f
|
||||
synchronized_at: 2026-09-05T07:16:44Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:14:27Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 开发工作流
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
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: 46067d78327476f14e6cfd202458a53f295bf763
|
||||
synchronized_at: 2026-09-19T07:18:28Z
|
||||
wiki_revision: 8e2cfa74bc7cf228278221e0f7ef488ac1b596d3
|
||||
synchronized_at: 2026-09-21T08:14:31Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 架构与代码地图
|
||||
@@ -486,6 +486,12 @@ Web 唯一展示位置为“采集采购 → SYB 同步记录”:列表状态
|
||||
- Web purchase-tasks API/页面:共享勾选、按动作筛选、独立状态列/详情及逐项接受结果;付款和物流流程保持原样。
|
||||
- 验证入口:go test ./app/goauto/purchase ./app/goauto/sybclient ./app/goauto/access ./app/goauto/migrations ./cmd/migrate/migration/version-local;Web tests/e2e/purchase-order-writeback.spec.ts。测试只使用隔离SQLite和fake/httptest,不代表MySQL多实例或真实SYB验收。
|
||||
|
||||
### 会话类失败自动重试(#330)
|
||||
|
||||
- purchase/order_writeback_worker.go:restoreOrderWritebackClient 在 ImportCookiesJSON 后调用 sybclient.CheckSession,UserID<=0 显式判不可用;finishSessionUnavailable 复用 lease_expires_at 作为退避到期时间(maxSessionRetryAttempts=6,sessionRetryBackoff 5/10/15/30/30m),领取条件增加 failed+SYB_SESSION_UNAVAILABLE+到期+未达上限。无迁移。
|
||||
- purchase/order_writeback.go:会话类失败的 CanSubmit 不受退避租约限制;手工重新提交 attempt_count 置 0。
|
||||
- 验证:go test ./app/goauto/purchase(含 httptest 模拟 /am/user/get 与断言 syb_session 未删除)。
|
||||
|
||||
## Chrome PDD 订单回填扩展(#316)
|
||||
|
||||
- `chrome-extension/` 是独立 Manifest V3 交付单元:popup 只负责配置、启动/停止、状态轮询和逐项结果;content script 只在 `mobile.yangkeduo.com` 的隔离世界中按可见 DOM 串行读取;service worker 负责持久运行状态、聚合冲突和分批提交。
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
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: b6df3e2d3d497827c9f3cc5b5ec859d5bda3c1f2
|
||||
synchronized_at: 2026-09-19T07:39:56Z
|
||||
wiki_revision: 6cb34c99ccdcf4b64b01cb9a25ff1549d8aeeec9
|
||||
synchronized_at: 2026-09-21T08:14:37Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 业务规则与术语
|
||||
@@ -332,7 +332,7 @@ synchronized_at: 2026-09-19T07:39:56Z
|
||||
- 一个 PDD 商品可能被多个虾皮商品共用,因此规格匹配按受影响的每个虾皮商品分项记录。主表只表达总体进度;Agent 展示和“继续采购”资格必须读取当前任务对应虾皮商品的分项状态。
|
||||
- 选错替代商品时,连续执行 B→C 不等于撤销 A→B,因为 B 可能还关联其他虾皮商品。正确纠错语义是把原 A→B 记录置为 `superseded`,再建立 A→C,并只处理原记录分项中冻结的影响集合。
|
||||
- 创建请求按 `create_request_id` 幂等;重放时源商品、替代商品、来源类型、来源任务、采集证据和发起设备必须一致,否则返回幂等冲突,不能静默覆盖。
|
||||
- `created_by_device_id` 只代表 Agent 设备。当前系统没有设备到采购员账号的绑定,多人多机场景若需要个人责任追踪,必须另建工单实现设备绑定操作员。
|
||||
- `created_by_device_id` 只代表 Agent 设备。设备可由管理员绑定到一个采购员账号;一个采购员可拥有多台设备,一台设备最多归属一个采购员,也允许暂不归属。SYB 商品一键关联/替换按当前采购员选择的归属设备读取该设备最新临时采集,不使用其他账号或其他设备的全局最新记录。
|
||||
|
||||
|
||||
## PDD 商品替换生效与规格匹配(#131)
|
||||
@@ -630,6 +630,10 @@ Web 唯一展示位置为“采集采购 → SYB 同步记录”:列表状态
|
||||
订单回填状态独立于采购成功、支付复核及物流 writeback_status。实付金额只存 Admin,SYB cost=0 为接口固定参数,不以金额推断付款。SYB 接口会同时更新采购状态/平台/时间,不能视作纯展示修改。
|
||||
|
||||
采购管理增加独立状态列、批量回填和详情补偿;复用既有访问权限,不增支付确认或审批。批量受理与最终成功分开展示;重试采购和回填分别筛选勾选项。远端无原子CAS,对系统外人工并发修改/超长延迟请求不能承诺绝对互斥;有冲突应人工核对,禁止强制覆盖。
|
||||
### SYB 会话类失败的有界自动重试(#330)
|
||||
|
||||
实现 01510a8/08b7095(2026-09-21,已合并 main,未部署、未生产验证)。回填 worker 从缓存会话恢复客户端后调用 SYB 会话校验;会话缺失/过期、串号失效、校验网络错误等均记为 `SYB_SESSION_UNAVAILABLE`,error_message 只记录类别和“将自动重试;如持续失败请恢复登录后重试”,不含原始错误。该类失败发生在任何写入之前,最多自动重试 6 次,退避 5/10/15/30/30 分钟(约 90 分钟,大于一个整点同步周期),达上限保持 failed 等人工。会话仍只由每小时 SYB 同步刷新;回填不登录、不 OCR、不删除或写入会话。其他失败码仍不自动重试。退避期内可手工重新回填,手工提交重置尝试次数。历史失败记录不会被自动领取。
|
||||
|
||||
## Agent 回填订单入口兼容(#307)
|
||||
|
||||
- 打开我的订单后最多采样8次,每次间隔1秒,连续两次订单导航结构一致才继续;登录/风控仍立即停止,超时给出明确原因。
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Local-Development-and-Verification
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Local-Development-and-Verification.-
|
||||
wiki_revision: f8996a09b2158253553c9e929eab72e1e53e8f61
|
||||
synchronized_at: 2026-09-18T07:57:06Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:14:42Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 本地开发与验证
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Common-Changes
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Common-Changes.-
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:30:12Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:14:50Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 常见修改指南
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Troubleshooting
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Troubleshooting
|
||||
wiki_revision: 18744477bfd17f396e8c76ec7a2fcc6720acdf6f
|
||||
synchronized_at: 2026-09-11T09:10:00Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:14:54Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 故障排查
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Product-Requirements-Overview
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Product-Requirements-Overview.-
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:30:20Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:00Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 产品需求总览与当前 MVP
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
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: 275caa306765440b0888dea701198d45481bfb40
|
||||
synchronized_at: 2026-09-19T07:40:19Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:05Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
<!-- gitea-wiki-mirror:start -->
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Delivery-Issues
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Delivery-Issues.-
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:30:28Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:09Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 当前 MVP 交付工单索引
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: OnePlus-Real-Device-Acceptance
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/OnePlus-Real-Device-Acceptance.-
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:30:32Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:13Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 一加真机验收记录
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: PDD-Detail-Rule-Migration-Analysis
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/PDD-Detail-Rule-Migration-Analysis.-
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:30:37Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:17Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# PDD 商品详情采集规则迁移分析
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: SYB-ERP-Interface-Contract
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/SYB-ERP-Interface-Contract.-
|
||||
wiki_revision: 432e392ebe428ea19c0aa260f8e5938f36e4c3f0
|
||||
synchronized_at: 2026-09-18T02:18:53Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:30Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 12 顺云宝(SYB)ERP 接口契约
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Deployment-and-Operations
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Deployment-and-Operations.-
|
||||
wiki_revision: c208aeeeb4baa3ae2da29f13f3651910c500e360
|
||||
synchronized_at: 2026-09-18T07:57:16Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:14:47Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 部署与运维
|
||||
|
||||
+2
-2
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Home
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Home
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:29:34Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:14:18Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# GoAuto 文档中心
|
||||
|
||||
Vendored
+2
-2
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Deployment-Template
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Deployment-Template.-
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:30:41Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:22Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 部署文档模板
|
||||
|
||||
Vendored
+2
-2
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Task-Archive-Template
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Task-Archive-Template.-
|
||||
wiki_revision: b1b1b343917e66288f4282bc6b3b90ea4ff3cca0
|
||||
synchronized_at: 2026-09-04T11:30:46Z
|
||||
wiki_revision: 03ea269058b50fea2be7842b9c284018987c82d9
|
||||
synchronized_at: 2026-09-21T08:15:26Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
> 本模板只用于用户明确要求的专项历史快照或读取既有归档,不属于标准任务闭环。单次任务的唯一事实来源是 Gitea 工单;不要为了完成普通任务创建本页面,也不要自动导出到 `docs/task/`。
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
|
||||
"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"
|
||||
)
|
||||
|
||||
@@ -29,6 +30,14 @@ func (handler Handler) List(context *gin.Context) {
|
||||
writeError(context, internalError(err))
|
||||
return
|
||||
}
|
||||
if role, _ := jwt.ExtractClaims(context)["rolekey"].(string); role != "admin" {
|
||||
id := currentUserID(context)
|
||||
if id > 0 {
|
||||
request.OwnerUserID = &id
|
||||
} else {
|
||||
request.OwnerUserID = new(uint64)
|
||||
}
|
||||
}
|
||||
response, err := NewService(db).List(context.Request.Context(), request)
|
||||
if err != nil {
|
||||
writeError(context, err)
|
||||
@@ -37,6 +46,68 @@ func (handler Handler) List(context *gin.Context) {
|
||||
context.JSON(http.StatusOK, gin.H{"code": http.StatusOK, "data": response})
|
||||
}
|
||||
|
||||
func currentUserID(c *gin.Context) uint64 {
|
||||
// #333: go-admin's Authorizator runs on every request with the
|
||||
// IdentityHandler map, which has no "user" entry, so c.Get("userId") is
|
||||
// always 0 there. The JWT "identity" claim is the authenticated user id.
|
||||
switch id := jwt.ExtractClaims(c)["identity"].(type) {
|
||||
case float64:
|
||||
if id > 0 {
|
||||
return uint64(id)
|
||||
}
|
||||
case int:
|
||||
if id > 0 {
|
||||
return uint64(id)
|
||||
}
|
||||
case int64:
|
||||
if id > 0 {
|
||||
return uint64(id)
|
||||
}
|
||||
case uint64:
|
||||
return id
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func (handler Handler) Owners(context *gin.Context) {
|
||||
db, err := handler.database(context)
|
||||
if err != nil {
|
||||
writeError(context, internalError(err))
|
||||
return
|
||||
}
|
||||
rows, err := NewService(db).Owners(context.Request.Context())
|
||||
if err != nil {
|
||||
writeError(context, err)
|
||||
return
|
||||
}
|
||||
context.JSON(http.StatusOK, gin.H{"code": http.StatusOK, "data": rows})
|
||||
}
|
||||
|
||||
func (handler Handler) SetOwner(context *gin.Context) {
|
||||
deviceID, err := strconv.ParseUint(context.Param("deviceId"), 10, 64)
|
||||
if err != nil || deviceID == 0 {
|
||||
writeError(context, invalidRequest("deviceId 无效"))
|
||||
return
|
||||
}
|
||||
var req struct {
|
||||
OwnerUserID *uint64 `json:"ownerUserId"`
|
||||
}
|
||||
if err := decodeJSON(context, &req); err != nil {
|
||||
writeError(context, invalidRequest("请求 JSON 无效"))
|
||||
return
|
||||
}
|
||||
db, err := handler.database(context)
|
||||
if err != nil {
|
||||
writeError(context, internalError(err))
|
||||
return
|
||||
}
|
||||
if err := NewService(db).SetOwner(context.Request.Context(), deviceID, req.OwnerUserID); err != nil {
|
||||
writeError(context, err)
|
||||
return
|
||||
}
|
||||
context.JSON(http.StatusOK, gin.H{"code": http.StatusOK, "data": gin.H{"deviceId": deviceID, "ownerUserId": req.OwnerUserID}})
|
||||
}
|
||||
|
||||
func (handler Handler) Register(context *gin.Context) {
|
||||
request, err := decodeRegisterRequest(context)
|
||||
if err != nil {
|
||||
|
||||
@@ -13,10 +13,11 @@ import (
|
||||
)
|
||||
|
||||
type ListRequest struct {
|
||||
Page int
|
||||
PageSize int
|
||||
Name string
|
||||
Status string
|
||||
Page int
|
||||
PageSize int
|
||||
Name string
|
||||
Status string
|
||||
OwnerUserID *uint64
|
||||
}
|
||||
|
||||
type DeviceListItem struct {
|
||||
@@ -36,6 +37,8 @@ type DeviceListItem struct {
|
||||
LastHeartbeatAt *time.Time `json:"lastHeartbeatAt"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
OwnerUserID *uint64 `json:"ownerUserId"`
|
||||
OwnerName string `json:"ownerName"`
|
||||
}
|
||||
|
||||
type DeviceListResponse struct {
|
||||
@@ -68,6 +71,9 @@ func (service *Service) List(ctx context.Context, request ListRequest) (DeviceLi
|
||||
if request.Status != "" {
|
||||
query = query.Where("status = ?", request.Status)
|
||||
}
|
||||
if request.OwnerUserID != nil {
|
||||
query = query.Where("owner_user_id = ?", *request.OwnerUserID)
|
||||
}
|
||||
var total int64
|
||||
if err := query.Count(&total).Error; err != nil {
|
||||
return DeviceListResponse{}, internalError(err)
|
||||
@@ -100,8 +106,8 @@ func (service *Service) List(ctx context.Context, request ListRequest) (DeviceLi
|
||||
AndroidVersion: record.AndroidVersion, AgentVersion: record.AgentVersion, PDDVersion: record.PDDVersion,
|
||||
Capabilities: capabilities,
|
||||
Status: record.Status, LastHeartbeatAt: record.LastHeartbeatAt,
|
||||
TokenRevoked: record.TokenRevokedAt != nil,
|
||||
CreatedAt: record.CreatedAt, UpdatedAt: record.UpdatedAt,
|
||||
TokenRevoked: record.TokenRevokedAt != nil, OwnerUserID: record.OwnerUserID,
|
||||
CreatedAt: record.CreatedAt, UpdatedAt: record.UpdatedAt,
|
||||
}
|
||||
if taskID, busy := currentTasks[record.ID]; busy {
|
||||
item.Busy = true
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
package device
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
|
||||
"go-admin/app/goauto/models"
|
||||
)
|
||||
|
||||
// #333: mirrors the real go-admin middleware, whose Authorizator sets
|
||||
// userId to 0 on every request while the JWT claims carry the identity.
|
||||
func listDevicesAs(t *testing.T, handler Handler, claims jwt.MapClaims) []DeviceListItem {
|
||||
t.Helper()
|
||||
gin.SetMode(gin.TestMode)
|
||||
engine := gin.New()
|
||||
engine.GET("/devices", func(c *gin.Context) {
|
||||
c.Set(jwt.JwtPayloadKey, claims)
|
||||
c.Set("userId", 0)
|
||||
c.Next()
|
||||
}, handler.List)
|
||||
recorder := httptest.NewRecorder()
|
||||
engine.ServeHTTP(recorder, httptest.NewRequest(http.MethodGet, "/devices?page=1&pageSize=20", nil))
|
||||
if recorder.Code != http.StatusOK {
|
||||
t.Fatalf("status=%d body=%s", recorder.Code, recorder.Body.String())
|
||||
}
|
||||
var body struct {
|
||||
Data DeviceListResponse `json:"data"`
|
||||
}
|
||||
if err := json.Unmarshal(recorder.Body.Bytes(), &body); err != nil {
|
||||
t.Fatalf("decode: %v", err)
|
||||
}
|
||||
return body.Data.Items
|
||||
}
|
||||
|
||||
func seedOwnedDevice(t *testing.T, handler Handler, name string, owner *uint64) models.AgentDevice {
|
||||
t.Helper()
|
||||
device := models.AgentDevice{
|
||||
InstallID: "install-" + name, Name: name, Manufacturer: "test", Model: "test",
|
||||
AndroidVersion: "14", AgentVersion: "1", PDDVersion: "1", CapabilitiesJSON: "[]",
|
||||
Status: models.DeviceStatusOnline, TokenDigest: fmt.Sprintf("digest-%s", name),
|
||||
TokenIssuedAt: time.Now(), OwnerUserID: owner,
|
||||
}
|
||||
if err := handler.DB.Create(&device).Error; err != nil {
|
||||
t.Fatalf("create device: %v", err)
|
||||
}
|
||||
return device
|
||||
}
|
||||
|
||||
func deviceNames(items []DeviceListItem) map[string]bool {
|
||||
names := map[string]bool{}
|
||||
for _, item := range items {
|
||||
names[item.Name] = true
|
||||
}
|
||||
return names
|
||||
}
|
||||
|
||||
func TestDeviceListPurchaserSeesOnlyOwnDevicesFromJWTIdentity(t *testing.T) {
|
||||
handler := Handler{DB: openTestDatabase(t)}
|
||||
two, three := uint64(2), uint64(3)
|
||||
seedOwnedDevice(t, handler, "caigou1-phone", &two)
|
||||
seedOwnedDevice(t, handler, "caigou2-phone", &three)
|
||||
seedOwnedDevice(t, handler, "unowned-phone", nil)
|
||||
|
||||
names := deviceNames(listDevicesAs(t, handler, jwt.MapClaims{"rolekey": "purchaser", "identity": float64(3)}))
|
||||
if len(names) != 1 || !names["caigou2-phone"] {
|
||||
t.Fatalf("purchaser 3 should see only own device, got %v", names)
|
||||
}
|
||||
|
||||
all := deviceNames(listDevicesAs(t, handler, jwt.MapClaims{"rolekey": "admin", "identity": float64(1)}))
|
||||
if len(all) != 3 {
|
||||
t.Fatalf("admin should see all devices, got %v", all)
|
||||
}
|
||||
|
||||
none := listDevicesAs(t, handler, jwt.MapClaims{"rolekey": "purchaser"})
|
||||
if len(none) != 0 {
|
||||
t.Fatalf("purchaser without identity must see no devices, got %d", len(none))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
package device
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"go-admin/app/goauto/models"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type DeviceOwnerItem struct {
|
||||
UserID uint64 `json:"userId"`
|
||||
Username string `json:"username"`
|
||||
NickName string `json:"nickName"`
|
||||
}
|
||||
|
||||
func (service *Service) Owners(ctx context.Context) ([]DeviceOwnerItem, error) {
|
||||
var rows []DeviceOwnerItem
|
||||
if err := service.DB.WithContext(ctx).Table("sys_user u").Select("u.user_id AS user_id, u.username, u.nick_name").Joins("JOIN sys_role r ON r.role_id = u.role_id").Where("u.status <> ? AND r.role_key <> ?", "1", "admin").Order("u.user_id").Scan(&rows).Error; err != nil {
|
||||
return nil, internalError(err)
|
||||
}
|
||||
return rows, nil
|
||||
}
|
||||
|
||||
func (service *Service) SetOwner(ctx context.Context, deviceID uint64, owner *uint64) error {
|
||||
if deviceID == 0 {
|
||||
return invalidRequest("deviceId 无效")
|
||||
}
|
||||
if owner != nil && *owner == 0 {
|
||||
return invalidRequest("ownerUserId 无效")
|
||||
}
|
||||
if owner != nil {
|
||||
var count int64
|
||||
if err := service.DB.WithContext(ctx).Table("sys_user u").Joins("JOIN sys_role r ON r.role_id = u.role_id").Where("u.user_id = ? AND u.status <> ? AND r.role_key <> ?", *owner, "1", "admin").Count(&count).Error; err != nil {
|
||||
return internalError(err)
|
||||
}
|
||||
if count != 1 {
|
||||
return invalidRequest("采购员不存在或不可分配")
|
||||
}
|
||||
}
|
||||
result := service.DB.WithContext(ctx).Model(&models.AgentDevice{}).Where("id = ?", deviceID).Update("owner_user_id", owner)
|
||||
if result.Error != nil {
|
||||
return internalError(result.Error)
|
||||
}
|
||||
if result.RowsAffected == 0 {
|
||||
return gorm.ErrRecordNotFound
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func parseOwnerID(value string) (*uint64, error) {
|
||||
value = strings.TrimSpace(value)
|
||||
if value == "" {
|
||||
return nil, nil
|
||||
}
|
||||
id, err := strconv.ParseUint(value, 10, 64)
|
||||
if err != nil || id == 0 {
|
||||
return nil, fmt.Errorf("ownerUserId 无效")
|
||||
}
|
||||
return &id, nil
|
||||
}
|
||||
@@ -21,6 +21,8 @@ func InitRouter(engine *gin.Engine, authMiddleware *jwt.GinJWTMiddleware) {
|
||||
|
||||
admin := engine.Group("/api/admin/v1/devices").Use(authMiddleware.MiddlewareFunc()).Use(middleware.AuthCheckRole())
|
||||
admin.GET("", handler.List)
|
||||
admin.GET("/owners", middleware.RequireRoleKey("admin"), handler.Owners)
|
||||
admin.PATCH("/:deviceId/owner", middleware.RequireRoleKey("admin"), handler.SetOwner)
|
||||
admin.POST("/:deviceId/disable", middleware.RequireRoleKey("admin"), handler.Disable)
|
||||
admin.POST("/:deviceId/identity-reset", middleware.RequireRoleKey("admin"), handler.ResetIdentity)
|
||||
admin.POST("/:deviceId/token/revoke", middleware.RequireRoleKey("admin"), handler.RevokeToken)
|
||||
|
||||
@@ -38,6 +38,7 @@ type AgentDevice struct {
|
||||
AndroidVersion string `json:"androidVersion" gorm:"size:32;not null"`
|
||||
AgentVersion string `json:"agentVersion" gorm:"size:32;not null"`
|
||||
PDDVersion string `json:"pddVersion" gorm:"size:32;not null"`
|
||||
OwnerUserID *uint64 `json:"ownerUserId" gorm:"column:owner_user_id;index"`
|
||||
CapabilitiesJSON string `json:"-" gorm:"size:4096;not null;default:'[]'"`
|
||||
Status string `json:"status" gorm:"size:16;not null;index;check:ck_agent_device_status,status IN ('online','offline','disabled')"`
|
||||
TokenDigest string `json:"-" gorm:"size:64;not null;uniqueIndex:ux_agent_device_token_digest"`
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
|
||||
"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"
|
||||
)
|
||||
|
||||
@@ -30,9 +31,16 @@ func (handler Handler) List(c *gin.Context) {
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
response, err := service.List(c.Request.Context(), ListRequest{
|
||||
request := ListRequest{
|
||||
Page: page, PageSize: pageSize, GoodsID: c.Query("goodsId"), Keyword: c.Query("keyword"), Status: strings.TrimSpace(c.Query("status")),
|
||||
})
|
||||
}
|
||||
claims := jwt.ExtractClaims(c)
|
||||
if role, _ := claims["rolekey"].(string); role != "admin" {
|
||||
id, _ := claims["identity"].(float64)
|
||||
owner := uint64(id)
|
||||
request.CollectionOwnerUserID = &owner
|
||||
}
|
||||
response, err := service.List(c.Request.Context(), request)
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
|
||||
@@ -64,10 +64,11 @@ type UpdateRequest struct {
|
||||
}
|
||||
|
||||
type ListRequest struct {
|
||||
Page, PageSize int
|
||||
Keyword string
|
||||
GoodsID string
|
||||
Status string
|
||||
Page, PageSize int
|
||||
Keyword string
|
||||
GoodsID string
|
||||
Status string
|
||||
CollectionOwnerUserID *uint64
|
||||
}
|
||||
|
||||
type ProductView struct {
|
||||
@@ -266,6 +267,9 @@ func (service *Service) List(ctx context.Context, request ListRequest) (ListResp
|
||||
}
|
||||
query = query.Where("status = ?", request.Status)
|
||||
}
|
||||
if request.CollectionOwnerUserID != nil {
|
||||
query = query.Where("EXISTS (SELECT 1 FROM collection_task ct JOIN agent_device ad ON ad.id = ct.device_id WHERE ct.pdd_product_id = pdd_product.id AND ct.source = ? AND ct.status IN ? AND ad.owner_user_id = ?)", models.CollectionTaskSourceAgentCurrentPage, []string{models.TaskStatusCompleted, models.TaskStatusCompletedPartial}, *request.CollectionOwnerUserID)
|
||||
}
|
||||
var total int64
|
||||
if err := query.Count(&total).Error; err != nil {
|
||||
return ListResponse{}, internalError(err)
|
||||
|
||||
@@ -205,6 +205,44 @@ func TestListMarksProductsUnavailableForCollection(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestListFiltersManualAssociationProductsByOwnedCollectionDevice(t *testing.T) {
|
||||
db := openProductDatabase(t)
|
||||
service := NewService(db)
|
||||
rule := models.CollectionRule{Name: "owned-rule", ContentJSON: `{}`}
|
||||
if err := db.Create(&rule).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
products := []models.PDDProduct{{GoodsID: "910001", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=910001", Status: "active"}, {GoodsID: "910002", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=910002", Status: "active"}, {GoodsID: "910003", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=910003", Status: "active"}}
|
||||
if err := db.Create(&products).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
owned := models.AgentDevice{InstallID: "owned-manual-list", Name: "owned", Manufacturer: "test", Model: "test", AndroidVersion: "14", AgentVersion: "1", PDDVersion: "1", Status: models.DeviceStatusOffline, TokenDigest: "owned-digest", TokenIssuedAt: time.Now(), OwnerUserID: ptrUint64(41)}
|
||||
other := models.AgentDevice{InstallID: "other-manual-list", Name: "other", Manufacturer: "test", Model: "test", AndroidVersion: "14", AgentVersion: "1", PDDVersion: "1", Status: models.DeviceStatusOffline, TokenDigest: "other-digest", TokenIssuedAt: time.Now(), OwnerUserID: ptrUint64(42)}
|
||||
if err := db.Create(&owned).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Create(&other).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, task := range []models.CollectionTask{{PDDProductID: &products[0].ID, DeviceID: &owned.ID, Source: models.CollectionTaskSourceAgentCurrentPage, Status: models.TaskStatusCompleted}, {PDDProductID: &products[1].ID, DeviceID: &other.ID, Source: models.CollectionTaskSourceAgentCurrentPage, Status: models.TaskStatusCompleted}, {PDDProductID: &products[2].ID, DeviceID: &owned.ID, Source: models.CollectionTaskSourceAdmin, Status: models.TaskStatusCompleted}} {
|
||||
task.RuleID = rule.ID
|
||||
task.URLSnapshot = products[0].URL
|
||||
task.RuleSnapshot = rule.ContentJSON
|
||||
if err := db.Create(&task).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
got, err := service.List(context.Background(), ListRequest{Page: 1, PageSize: 20, CollectionOwnerUserID: ptrUint64(41)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(got.Items) != 1 || got.Items[0].GoodsID != products[0].GoodsID {
|
||||
t.Fatalf("unexpected owned products: %+v", got.Items)
|
||||
}
|
||||
}
|
||||
|
||||
func ptrUint64(value uint64) *uint64 { return &value }
|
||||
|
||||
func TestListDisablesCollectionWhenNoRuleExists(t *testing.T) {
|
||||
db := openProductDatabase(t)
|
||||
product := models.PDDProduct{GoodsID: "444444", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=444444", Status: "active"}
|
||||
|
||||
@@ -108,7 +108,12 @@ func (s *Service) OrderWritebackViews(ctx context.Context, tasks []models.Purcha
|
||||
v.Reason = r.ErrorMessage
|
||||
}
|
||||
v.CompletedAt = r.CompletedAt
|
||||
v.CanSubmit = v.CanSubmit && (r.Status == "failed" || r.Status == "unknown") && (r.LeaseExpiresAt == nil || !r.LeaseExpiresAt.After(s.Now()))
|
||||
// A session-class failure's LeaseExpiresAt is the automatic-retry backoff
|
||||
// deadline (#330 修订2, order_writeback_worker.go finishSessionUnavailable),
|
||||
// not an in-flight write lease — manual "resubmit" must stay available
|
||||
// during that window instead of being hidden until it expires.
|
||||
sessionBackoff := r.Status == "failed" && r.ErrorCode == "SYB_SESSION_UNAVAILABLE"
|
||||
v.CanSubmit = v.CanSubmit && (r.Status == "failed" || r.Status == "unknown") && (sessionBackoff || r.LeaseExpiresAt == nil || !r.LeaseExpiresAt.After(s.Now()))
|
||||
out[r.PurchaseTaskID] = v
|
||||
}
|
||||
return out, nil
|
||||
@@ -197,7 +202,12 @@ func (s *Service) RequestOrderWriteback(ctx context.Context, req OrderWritebackR
|
||||
case row.Status == "pending":
|
||||
a.Result, a.Reason = "pending", "已加入回填"
|
||||
default:
|
||||
if err := tx.Model(&row).Updates(map[string]any{"status": "pending", "write_started": false, "lease_owner": "", "lease_expires_at": nil, "error_code": "", "error_message": ""}).Error; err != nil {
|
||||
// Manual resubmit resets attempt_count to 0 (#330 修订3) so a stale
|
||||
// history of automatic session-class retries never eats into a
|
||||
// fresh manual attempt budget, and clears lease_expires_at so the
|
||||
// worker's bounded auto-retry claim (which requires it non-nil)
|
||||
// does not race a double-claim against this manual pending row.
|
||||
if err := tx.Model(&row).Updates(map[string]any{"status": "pending", "write_started": false, "lease_owner": "", "lease_expires_at": nil, "error_code": "", "error_message": "", "attempt_count": 0}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
a.Result, a.Reason = "pending", "已加入回填,将先回读SYB"
|
||||
|
||||
@@ -0,0 +1,391 @@
|
||||
package purchase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/app/goauto/sybclient"
|
||||
"go-admin/config"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// wbFactoryWorker builds a worker whose Factory itself fails, exercising the
|
||||
// restoreOrderWritebackClient failure path (session unavailable) rather than
|
||||
// a remote read/write failure on an otherwise-working client.
|
||||
func wbFactoryWorker(s *Service, err error) *OrderWritebackWorker {
|
||||
return &OrderWritebackWorker{DB: s.DB, Now: s.Now, Factory: func(context.Context, *gorm.DB) (OrderNumberClient, error) { return nil, err }}
|
||||
}
|
||||
|
||||
func TestOrderWritebackSessionFailureSchedulesBoundedRetry(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
if ok, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil || !ok {
|
||||
t.Fatalf("run %v %v", ok, err)
|
||||
}
|
||||
row := loadOrderWriteback(t, s.DB, task.ID)
|
||||
if row.Status != "failed" || row.ErrorCode != "SYB_SESSION_UNAVAILABLE" {
|
||||
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||
}
|
||||
if row.ErrorMessage == "" || len(row.ErrorMessage) > 300 {
|
||||
t.Fatalf("error message not recorded safely: %q", row.ErrorMessage)
|
||||
}
|
||||
if !strings.Contains(row.ErrorMessage, "会话缺失/已过期") {
|
||||
t.Fatalf("category missing from message: %q", row.ErrorMessage)
|
||||
}
|
||||
if !strings.Contains(row.ErrorMessage, "将自动重试") || !strings.Contains(row.ErrorMessage, "恢复登录") {
|
||||
t.Fatalf("message is not actionable: %q", row.ErrorMessage)
|
||||
}
|
||||
if row.LeaseExpiresAt == nil || !row.LeaseExpiresAt.After(s.Now()) {
|
||||
t.Fatal("no backoff scheduled for first session-class failure")
|
||||
}
|
||||
if row.AttemptCount != 1 {
|
||||
t.Fatalf("attempt_count=%d", row.AttemptCount)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSessionRetryBackoffTotalExceedsHourlySyncWindow guards the ticket's
|
||||
// blocker: the cumulative auto-retry window must outlast one hourly sync
|
||||
// period (up to ~60 minutes from failure to the refresh that fixes it),
|
||||
// otherwise attempts run out before the session has a chance to recover.
|
||||
func TestSessionRetryBackoffTotalExceedsHourlySyncWindow(t *testing.T) {
|
||||
if len(sessionRetryBackoff) != maxSessionRetryAttempts-1 {
|
||||
t.Fatalf("expected %d backoff steps for %d attempts, got %d", maxSessionRetryAttempts-1, maxSessionRetryAttempts, len(sessionRetryBackoff))
|
||||
}
|
||||
var total time.Duration
|
||||
for _, d := range sessionRetryBackoff {
|
||||
total += d
|
||||
}
|
||||
if total <= time.Hour {
|
||||
t.Fatalf("total backoff %s must exceed one hourly sync period", total)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderWritebackSessionFailureNotReclaimedBeforeBackoffExpires(t *testing.T) {
|
||||
s, _ := orderWritebackFixture(t)
|
||||
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
f := &fakeOrderNumberClient{apply: true}
|
||||
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || ok {
|
||||
t.Fatalf("claimed before backoff expired: ok=%v err=%v", ok, err)
|
||||
}
|
||||
if f.writes != 0 {
|
||||
t.Fatal("wrote while still inside backoff window")
|
||||
}
|
||||
}
|
||||
|
||||
// loadOrderWriteback in order_writeback_test.go takes (t, db, id); provide a
|
||||
// small adapter so this file reads naturally when task id is already in hand.
|
||||
func loadOrderWritebackByTask(t *testing.T, s *Service, id uint64) models.PurchaseOrderWriteback {
|
||||
return loadOrderWriteback(t, s.DB, id)
|
||||
}
|
||||
|
||||
func TestOrderWritebackSessionFailureReclaimedAfterBackoffExpires(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||
s.Now = func() time.Time { return row.LeaseExpiresAt.Add(time.Second) }
|
||||
f := &fakeOrderNumberClient{apply: true}
|
||||
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || !ok {
|
||||
t.Fatalf("not reclaimed after backoff expired: ok=%v err=%v", ok, err)
|
||||
}
|
||||
if f.writes != 1 {
|
||||
t.Fatal("did not write after successful reclaim")
|
||||
}
|
||||
after := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if after.Status != "succeeded" {
|
||||
t.Fatalf("status=%s", after.Status)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderWritebackSessionFailureStopsRetryingAtMaxAttempts(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
now := s.Now()
|
||||
for i := 0; i < maxSessionRetryAttempts; i++ {
|
||||
s.Now = func() time.Time { return now }
|
||||
if ok, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil || !ok {
|
||||
t.Fatalf("attempt %d: ok=%v err=%v", i+1, ok, err)
|
||||
}
|
||||
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if row.AttemptCount != i+1 {
|
||||
t.Fatalf("attempt %d: attempt_count=%d", i+1, row.AttemptCount)
|
||||
}
|
||||
if i+1 < maxSessionRetryAttempts {
|
||||
if row.LeaseExpiresAt == nil {
|
||||
t.Fatalf("attempt %d: no backoff scheduled", i+1)
|
||||
}
|
||||
now = row.LeaseExpiresAt.Add(time.Second)
|
||||
} else {
|
||||
if row.LeaseExpiresAt != nil {
|
||||
t.Fatal("lease still scheduled at max attempts")
|
||||
}
|
||||
}
|
||||
}
|
||||
// One more tick past any plausible backoff: the claim query must exclude
|
||||
// attempt_count >= maxSessionRetryAttempts, so nothing is claimed.
|
||||
s.Now = func() time.Time { return now.Add(24 * time.Hour) }
|
||||
f := &fakeOrderNumberClient{apply: true}
|
||||
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || ok {
|
||||
t.Fatalf("claimed a row past max attempts: ok=%v err=%v", ok, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderWritebackCheckSessionInvalidIsSessionClassAndNeverDeletesSession(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
seedWritebackSession(t, s)
|
||||
if ok, err := wbFactoryWorker(s, sybclient.ErrSessionInvalid).RunOnce(context.Background()); err != nil || !ok {
|
||||
t.Fatalf("run %v %v", ok, err)
|
||||
}
|
||||
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if row.Status != "failed" || row.ErrorCode != "SYB_SESSION_UNAVAILABLE" {
|
||||
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||
}
|
||||
if row.LeaseExpiresAt == nil {
|
||||
t.Fatal("ErrSessionInvalid was not scheduled for retry")
|
||||
}
|
||||
assertWritebackSessionUntouched(t, s)
|
||||
}
|
||||
|
||||
func TestOrderWritebackCheckSessionNetworkErrorIsSessionClassAndNeverDeletesSession(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
seedWritebackSession(t, s)
|
||||
if ok, err := wbFactoryWorker(s, errors.New("dial tcp: i/o timeout")).RunOnce(context.Background()); err != nil || !ok {
|
||||
t.Fatalf("run %v %v", ok, err)
|
||||
}
|
||||
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if row.Status != "failed" || row.ErrorCode != "SYB_SESSION_UNAVAILABLE" {
|
||||
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||
}
|
||||
if row.LeaseExpiresAt == nil {
|
||||
t.Fatal("network error was not scheduled for retry")
|
||||
}
|
||||
assertWritebackSessionUntouched(t, s)
|
||||
}
|
||||
|
||||
func seedWritebackSession(t *testing.T, s *Service) {
|
||||
t.Helper()
|
||||
store := sybclient.NewSessionStore(s.DB)
|
||||
// restoreOrderWritebackClient's SessionStore.Load compares against real
|
||||
// wall-clock time.Now(), not the service's mocked s.Now (which fixtures
|
||||
// pin to a fixed past date) — so the session must expire relative to the
|
||||
// real clock or Load reports ErrNoSession even though a row exists.
|
||||
if err := store.Save(context.Background(), sybclient.Session{
|
||||
Username: "syb-writeback-test", UserID: 555, CookiesJSON: `[{"name":"SESSION","value":"x"}]`, ExpiresAt: time.Now().Add(time.Hour),
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func assertWritebackSessionUntouched(t *testing.T, s *Service) {
|
||||
t.Helper()
|
||||
var count int64
|
||||
if err := s.DB.Model(&models.SYBSession{}).Where("username = ?", "syb-writeback-test").Count(&count).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Fatal("writeback worker deleted or otherwise removed the cached SYB session")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderWritebackCanSubmitDuringSessionBackoff(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
updated := loadBackfillTask(t, s.DB, task.ID)
|
||||
views, err := s.OrderWritebackViews(context.Background(), []models.PurchaseTask{updated})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
v := views[task.ID]
|
||||
if v.Status != "failed" || !v.CanSubmit {
|
||||
t.Fatalf("expected resubmit available during backoff: %+v", v)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderWritebackManualResubmitResetsAttemptCountAndLease(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
now := s.Now()
|
||||
for i := 0; i < 3; i++ {
|
||||
s.Now = func() time.Time { return now }
|
||||
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||
now = row.LeaseExpiresAt.Add(time.Second)
|
||||
}
|
||||
before := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if before.AttemptCount != 3 {
|
||||
t.Fatalf("attempt_count=%d", before.AttemptCount)
|
||||
}
|
||||
if _, err := s.RequestOrderWriteback(context.Background(), OrderWritebackRequest{uuid.NewString(), []uint64{task.ID}}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
after := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if after.Status != "pending" || after.AttemptCount != 0 || after.LeaseExpiresAt != nil {
|
||||
t.Fatalf("manual resubmit did not reset state: %+v", after)
|
||||
}
|
||||
f := &fakeOrderNumberClient{apply: true}
|
||||
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || !ok || f.writes != 1 {
|
||||
t.Fatalf("worker could not process post-resubmit row: ok=%v err=%v writes=%d", ok, err, f.writes)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOrderWritebackOtherFailureCodesAreNotAutoRetried(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
f := &fakeOrderNumberClient{readErr: errors.New("offline")}
|
||||
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || !ok {
|
||||
t.Fatalf("run %v %v", ok, err)
|
||||
}
|
||||
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if row.Status != "failed" || row.ErrorCode != "SYB_READ_FAILED" {
|
||||
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||
}
|
||||
if row.LeaseExpiresAt != nil {
|
||||
t.Fatal("non-session failure code must not be scheduled for automatic retry")
|
||||
}
|
||||
s.Now = func() time.Time { return row.CreatedAt.Add(24 * time.Hour) }
|
||||
again := &fakeOrderNumberClient{apply: true}
|
||||
if ok, err := wbWorker(s, again).RunOnce(context.Background()); err != nil || ok {
|
||||
t.Fatalf("a non-session failure code was auto-reclaimed: ok=%v err=%v", ok, err)
|
||||
}
|
||||
}
|
||||
|
||||
// --- restoreOrderWritebackClient against a real sybclient.Client + emulated
|
||||
// SYB /am/user/get, so the CheckSession probe added by #330 is actually
|
||||
// exercised end to end instead of only through a fake Factory. ---
|
||||
|
||||
// sybUserGetServer emulates the one endpoint restoreOrderWritebackClient's
|
||||
// CheckSession call depends on, using the real envelope shape documented in
|
||||
// sybclient/client.go's `envelope` type and asserted against in
|
||||
// sybclient/client_test.go.
|
||||
func sybUserGetServer(t *testing.T, handler func(w http.ResponseWriter, r *http.Request)) *httptest.Server {
|
||||
t.Helper()
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/am/user/get", handler)
|
||||
srv := httptest.NewServer(mux)
|
||||
t.Cleanup(srv.Close)
|
||||
return srv
|
||||
}
|
||||
|
||||
// withWritebackSYBConfig points config.ExtConfig.SYB at the given test
|
||||
// server for the duration of the test, restoring the previous value
|
||||
// afterwards so other tests (and any parallel config reads) are unaffected.
|
||||
func withWritebackSYBConfig(t *testing.T, baseURL, username string) {
|
||||
t.Helper()
|
||||
prev := config.ExtConfig.SYB
|
||||
config.ExtConfig.SYB = config.SYB{BaseURL: baseURL, Username: username}
|
||||
t.Cleanup(func() { config.ExtConfig.SYB = prev })
|
||||
}
|
||||
|
||||
func envelopeOK(w http.ResponseWriter, data any) {
|
||||
body, _ := json.Marshal(data)
|
||||
env, _ := json.Marshal(map[string]any{"status": true, "msg": "获取成功", "data": json.RawMessage(body), "code": nil})
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.Write(env)
|
||||
}
|
||||
|
||||
func TestRestoreOrderWritebackClientValidSessionReturnsClient(t *testing.T) {
|
||||
s, _ := orderWritebackFixture(t)
|
||||
seedWritebackSession(t, s)
|
||||
srv := sybUserGetServer(t, func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Query().Get("id") != strconv.FormatInt(555, 10) {
|
||||
t.Fatalf("unexpected id query: %s", r.URL.RawQuery)
|
||||
}
|
||||
envelopeOK(w, map[string]any{"id": 555, "username": "syb-writeback-test"})
|
||||
})
|
||||
withWritebackSYBConfig(t, srv.URL, "syb-writeback-test")
|
||||
client, err := restoreOrderWritebackClient(context.Background(), s.DB)
|
||||
if err != nil || client == nil {
|
||||
t.Fatalf("expected a usable client, got client=%v err=%v", client, err)
|
||||
}
|
||||
assertWritebackSessionUntouched(t, s)
|
||||
}
|
||||
|
||||
func TestRestoreOrderWritebackClientMismatchedUsernameIsSessionInvalid(t *testing.T) {
|
||||
s, _ := orderWritebackFixture(t)
|
||||
seedWritebackSession(t, s)
|
||||
srv := sybUserGetServer(t, func(w http.ResponseWriter, r *http.Request) {
|
||||
// SYB says the cookie now belongs to a different account (12 §3.5):
|
||||
// treated the same as an explicit logout.
|
||||
envelopeOK(w, map[string]any{"id": 555, "username": "somebody-else"})
|
||||
})
|
||||
withWritebackSYBConfig(t, srv.URL, "syb-writeback-test")
|
||||
_, err := restoreOrderWritebackClient(context.Background(), s.DB)
|
||||
if !errors.Is(err, sybclient.ErrSessionInvalid) {
|
||||
t.Fatalf("expected ErrSessionInvalid, got %v", err)
|
||||
}
|
||||
assertWritebackSessionUntouched(t, s)
|
||||
}
|
||||
|
||||
func TestRestoreOrderWritebackClient500IsNotSessionInvalid(t *testing.T) {
|
||||
s, _ := orderWritebackFixture(t)
|
||||
seedWritebackSession(t, s)
|
||||
srv := sybUserGetServer(t, func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
})
|
||||
withWritebackSYBConfig(t, srv.URL, "syb-writeback-test")
|
||||
_, err := restoreOrderWritebackClient(context.Background(), s.DB)
|
||||
if err == nil {
|
||||
t.Fatal("expected an error for a 5xx response")
|
||||
}
|
||||
if errors.Is(err, sybclient.ErrSessionInvalid) {
|
||||
t.Fatalf("a 5xx must not be classified as a confirmed logout, got %v", err)
|
||||
}
|
||||
assertWritebackSessionUntouched(t, s)
|
||||
}
|
||||
|
||||
func TestRestoreOrderWritebackClientTimeoutIsNotSessionInvalid(t *testing.T) {
|
||||
s, _ := orderWritebackFixture(t)
|
||||
seedWritebackSession(t, s)
|
||||
srv := sybUserGetServer(t, func(w http.ResponseWriter, r *http.Request) {
|
||||
<-r.Context().Done() // never respond; the client-side ctx timeout fires first
|
||||
})
|
||||
withWritebackSYBConfig(t, srv.URL, "syb-writeback-test")
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond)
|
||||
defer cancel()
|
||||
_, err := restoreOrderWritebackClient(ctx, s.DB)
|
||||
if err == nil {
|
||||
t.Fatal("expected an error for a request that never completes")
|
||||
}
|
||||
if errors.Is(err, sybclient.ErrSessionInvalid) {
|
||||
t.Fatalf("a timeout must not be classified as a confirmed logout, got %v", err)
|
||||
}
|
||||
assertWritebackSessionUntouched(t, s)
|
||||
}
|
||||
|
||||
// TestOrderWritebackWorkerWithRealFactoryOnInvalidSession is the requested
|
||||
// end-to-end case: the worker's actual Factory (restoreOrderWritebackClient)
|
||||
// against a server that reports the cached session invalid. It must record
|
||||
// SYB_SESSION_UNAVAILABLE with a scheduled backoff and must not touch the
|
||||
// cached session row.
|
||||
func TestOrderWritebackWorkerWithRealFactoryOnInvalidSession(t *testing.T) {
|
||||
s, task := orderWritebackFixture(t)
|
||||
seedWritebackSession(t, s)
|
||||
srv := sybUserGetServer(t, func(w http.ResponseWriter, r *http.Request) {
|
||||
envelopeOK(w, map[string]any{"id": 555, "username": "somebody-else"})
|
||||
})
|
||||
withWritebackSYBConfig(t, srv.URL, "syb-writeback-test")
|
||||
w := &OrderWritebackWorker{DB: s.DB, Now: s.Now, Factory: restoreOrderWritebackClient}
|
||||
if ok, err := w.RunOnce(context.Background()); err != nil || !ok {
|
||||
t.Fatalf("run %v %v", ok, err)
|
||||
}
|
||||
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||
if row.Status != "failed" || row.ErrorCode != "SYB_SESSION_UNAVAILABLE" {
|
||||
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||
}
|
||||
if row.LeaseExpiresAt == nil {
|
||||
t.Fatal("no backoff scheduled")
|
||||
}
|
||||
assertWritebackSessionUntouched(t, s)
|
||||
}
|
||||
@@ -24,12 +24,80 @@ type OrderWritebackWorker struct {
|
||||
Factory func(context.Context, *gorm.DB) (OrderNumberClient, error)
|
||||
}
|
||||
|
||||
// errSessionUserIDMissing marks a cached session whose UserID column is not a
|
||||
// positive SYB account id. SessionStore.Save (session.go) rejects UserID<=0
|
||||
// before it is ever persisted, so this should be unreachable in practice; it
|
||||
// exists so a corrupted/legacy row fails loudly and safely instead of calling
|
||||
// CheckSession with id=0 (#330 修订1).
|
||||
var errSessionUserIDMissing = errors.New("SYB 会话记录缺少有效 user id")
|
||||
|
||||
// Bounded auto-retry for session-class writeback failures (#330). A session
|
||||
// outage self-heals once GoAutoSYBHourlySync refreshes syb_session, but that
|
||||
// refresh only happens once per hour (at :05) and only fires the run *after*
|
||||
// the session is found dead — so the wait from failure to refresh can be
|
||||
// close to a full hour. The backoff schedule below sums to ~90 minutes
|
||||
// (5+10+15+30+30) across maxSessionRetryAttempts=6 attempts, deliberately
|
||||
// longer than one hourly sync period so a session recovered by "the next"
|
||||
// hourly run is still caught automatically instead of exhausting attempts
|
||||
// first. maxSessionRetryAttempts caps the automatic attempts so a session
|
||||
// that never recovers still lands back in "failed" for a human instead of
|
||||
// retrying forever.
|
||||
const maxSessionRetryAttempts = 6
|
||||
|
||||
var sessionRetryBackoff = []time.Duration{
|
||||
5 * time.Minute,
|
||||
10 * time.Minute,
|
||||
15 * time.Minute,
|
||||
30 * time.Minute,
|
||||
30 * time.Minute,
|
||||
}
|
||||
|
||||
// sessionRetryDelay returns the backoff before the next automatic attempt,
|
||||
// given the attempt number (1-based, i.e. the count already recorded for the
|
||||
// attempt that just failed).
|
||||
func sessionRetryDelay(attempt int) time.Duration {
|
||||
idx := attempt - 1
|
||||
if idx < 0 {
|
||||
idx = 0
|
||||
}
|
||||
if idx >= len(sessionRetryBackoff) {
|
||||
idx = len(sessionRetryBackoff) - 1
|
||||
}
|
||||
return sessionRetryBackoff[idx]
|
||||
}
|
||||
|
||||
// sessionUnavailableMessage classifies why the cached SYB session could not
|
||||
// be used, without ever including cookies, tokens or other credential
|
||||
// material (#330 修订1点3). The category — not the raw error text — is what
|
||||
// gets persisted to error_message, wrapped in a fixed, actionable template
|
||||
// that stays well under the 300-char column limit.
|
||||
func sessionUnavailableMessage(err error) string {
|
||||
category := "会话恢复失败(网络/其他)"
|
||||
switch {
|
||||
case errors.Is(err, sybclient.ErrNoSession):
|
||||
category = "会话缺失/已过期"
|
||||
case errors.Is(err, errSessionUserIDMissing):
|
||||
category = "会话记录异常,缺少 user id"
|
||||
case errors.Is(err, sybclient.ErrSessionInvalid):
|
||||
category = "会话校验失效"
|
||||
}
|
||||
return "SYB会话不可用(" + category + "),将自动重试;如持续失败请恢复登录后重试"
|
||||
}
|
||||
|
||||
// restoreOrderWritebackClient rebuilds a SYB client from the cached session
|
||||
// only. It never logs in, never triggers OCR and never deletes the cached
|
||||
// session (that stays the exclusive responsibility of sybimport.Connect's
|
||||
// login/refresh path) — it only reports whether the cached cookies still
|
||||
// work, via CheckSession, so the caller can classify the failure (#330).
|
||||
func restoreOrderWritebackClient(ctx context.Context, db *gorm.DB) (OrderNumberClient, error) {
|
||||
cfg := config.ExtConfig.SYB.Resolved()
|
||||
session, err := sybclient.NewSessionStore(db).Load(ctx, cfg.Username, time.Now())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if session.UserID <= 0 {
|
||||
return nil, errSessionUserIDMissing
|
||||
}
|
||||
c, err := sybclient.New(cfg.BaseURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -37,6 +105,16 @@ func restoreOrderWritebackClient(ctx context.Context, db *gorm.DB) (OrderNumberC
|
||||
if err = c.ImportCookiesJSON(session.CookiesJSON); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// Active probe (#330 修订1): without this, a remotely-expired cookie jar
|
||||
// imports cleanly and only fails later inside read(), which would record
|
||||
// it as SYB_READ_FAILED instead of the retryable session-class outcome.
|
||||
// Any error here — ErrSessionInvalid or network/format — is treated as
|
||||
// session-class; only ErrSessionInvalid is a confirmed logout, but a
|
||||
// network/format error is not confirmed-valid either, so it is still
|
||||
// retried rather than attempted as a write.
|
||||
if err = c.CheckSession(ctx, session.UserID, cfg.Username); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return c, nil
|
||||
}
|
||||
|
||||
@@ -78,7 +156,10 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
||||
var item models.PurchaseOrderWriteback
|
||||
recovering := false
|
||||
err := db.Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("status = ? OR (status = ? AND lease_expires_at <= ?)", "pending", "running", now).Order("id").First(&item).Error; err != nil {
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where(
|
||||
"status = ? OR (status = ? AND lease_expires_at <= ?) OR (status = ? AND error_code = ? AND lease_expires_at IS NOT NULL AND lease_expires_at <= ? AND attempt_count < ?)",
|
||||
"pending", "running", now, "failed", "SYB_SESSION_UNAVAILABLE", now, maxSessionRetryAttempts,
|
||||
).Order("id").First(&item).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
recovering = item.Status == "running"
|
||||
@@ -90,6 +171,11 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
// tx.Model(&item).Updates used gorm.Expr("attempt_count + 1") above, which
|
||||
// GORM does not read back into the struct; sync it here so downstream
|
||||
// bounded-retry math (finishSessionUnavailable) sees the true post-claim
|
||||
// count instead of being off by one.
|
||||
item.AttemptCount++
|
||||
finish := func(status, code, message string) error {
|
||||
updates := map[string]any{"status": status, "error_code": code, "error_message": message, "lease_owner": ""}
|
||||
if status != "unknown" {
|
||||
@@ -100,6 +186,23 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
||||
}
|
||||
return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error
|
||||
}
|
||||
// finishSessionUnavailable is the bounded-retry counterpart of finish for
|
||||
// SYB_SESSION_UNAVAILABLE: instead of clearing the lease, it schedules the
|
||||
// next automatic attempt (item.AttemptCount was already incremented by the
|
||||
// claim above) until maxSessionRetryAttempts is reached, at which point it
|
||||
// behaves like finish("failed", ...) and stops retrying (#330).
|
||||
finishSessionUnavailable := func(err error) error {
|
||||
updates := map[string]any{
|
||||
"status": "failed", "error_code": "SYB_SESSION_UNAVAILABLE",
|
||||
"error_message": sessionUnavailableMessage(err), "lease_owner": "",
|
||||
}
|
||||
if item.AttemptCount < maxSessionRetryAttempts {
|
||||
updates["lease_expires_at"] = w.Now().Add(sessionRetryDelay(item.AttemptCount))
|
||||
} else {
|
||||
updates["lease_expires_at"] = nil
|
||||
}
|
||||
return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error
|
||||
}
|
||||
var task models.PurchaseTask
|
||||
if err = db.First(&task, item.PurchaseTaskID).Error; err != nil {
|
||||
return true, finish("failed", "TASK_UNAVAILABLE", "采购任务不可用,请人工核对")
|
||||
@@ -115,7 +218,7 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
||||
client, err := w.Factory(callCtx, db)
|
||||
cancel()
|
||||
if err != nil {
|
||||
return true, finish("failed", "SYB_SESSION_UNAVAILABLE", "SYB会话不可用,请恢复登录后重试")
|
||||
return true, finishSessionUnavailable(err)
|
||||
}
|
||||
read := func() (string, string, error) {
|
||||
readCtx, stop := context.WithTimeout(ctx, 20*time.Second)
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
package shopeeproduct
|
||||
|
||||
import (
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
|
||||
)
|
||||
|
||||
// #333: the real middleware sets userId=0 on every request; the operator used
|
||||
// for owned-device collection lookup and audit must come from the JWT identity.
|
||||
func TestCurrentUserIDReadsJWTIdentityNotUserIDKey(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
cases := []struct {
|
||||
claims jwt.MapClaims
|
||||
want uint64
|
||||
}{
|
||||
{jwt.MapClaims{"rolekey": "purchaser", "identity": float64(3)}, 3},
|
||||
{jwt.MapClaims{"rolekey": "purchaser"}, 0},
|
||||
{jwt.MapClaims{"rolekey": "purchaser", "identity": float64(-1)}, 0},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
c, _ := gin.CreateTestContext(httptest.NewRecorder())
|
||||
c.Set(jwt.JwtPayloadKey, tc.claims)
|
||||
c.Set("userId", 0)
|
||||
if got := currentUserID(c); got != tc.want {
|
||||
t.Fatalf("claims %v: got %d want %d", tc.claims, got, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
|
||||
"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"
|
||||
)
|
||||
|
||||
@@ -120,6 +121,8 @@ func (handler Handler) LinkPDD(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
service.ReplacementEligibility = handler.ReplacementEligibility
|
||||
service.OperatorUserID = currentUserID(c)
|
||||
service.OperatorIsAdmin = currentRole(c) == "admin"
|
||||
response, err := service.LinkPDD(c.Request.Context(), id, request)
|
||||
respond(c, response, err)
|
||||
}
|
||||
@@ -424,22 +427,31 @@ func (handler Handler) service(c *gin.Context) (*Service, bool) {
|
||||
// purchasing roles per #40; role membership itself is enforced by the router
|
||||
// middleware, not here.
|
||||
func currentUserID(c *gin.Context) uint64 {
|
||||
value, exists := c.Get("userId")
|
||||
if !exists {
|
||||
return 0
|
||||
}
|
||||
switch id := value.(type) {
|
||||
// #333: go-admin's Authorizator runs on every request with the
|
||||
// IdentityHandler map, which has no "user" entry, so c.Get("userId") is
|
||||
// always 0 there. The JWT "identity" claim is the authenticated user id.
|
||||
switch id := jwt.ExtractClaims(c)["identity"].(type) {
|
||||
case float64:
|
||||
if id > 0 {
|
||||
return uint64(id)
|
||||
}
|
||||
case int:
|
||||
return uint64(id)
|
||||
if id > 0 {
|
||||
return uint64(id)
|
||||
}
|
||||
case int64:
|
||||
return uint64(id)
|
||||
if id > 0 {
|
||||
return uint64(id)
|
||||
}
|
||||
case uint64:
|
||||
return id
|
||||
case float64:
|
||||
return uint64(id)
|
||||
default:
|
||||
return 0
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func currentRole(c *gin.Context) string {
|
||||
value, _ := jwt.ExtractClaims(c)["rolekey"].(string)
|
||||
return value
|
||||
}
|
||||
|
||||
func respond(c *gin.Context, response SaveResponse, err error) {
|
||||
@@ -489,6 +501,8 @@ func writeError(c *gin.Context, err error) {
|
||||
status = http.StatusNotFound
|
||||
case CodePDDProductDisabled, CodeSpecContextStale, CodeLatestCollectionUnavailable, CodeLinkConflict:
|
||||
status = http.StatusConflict
|
||||
case CodeDeviceOwnershipForbidden:
|
||||
status = http.StatusForbidden
|
||||
case CodeAIUnavailable:
|
||||
status = http.StatusServiceUnavailable
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
|
||||
const CodeLatestCollectionUnavailable = "LATEST_COLLECTION_UNAVAILABLE"
|
||||
const CodeLinkConflict = "PDD_LINK_CONFLICT"
|
||||
const CodeDeviceOwnershipForbidden = "DEVICE_OWNERSHIP_FORBIDDEN"
|
||||
|
||||
type LatestCollectionRequest struct {
|
||||
SYBProductID uint64 `json:"sybProductId"`
|
||||
@@ -105,6 +106,9 @@ func (service *Service) linkLatestCollection(ctx context.Context, id uint64, req
|
||||
if device.Status != models.DeviceStatusOnline && device.Status != models.DeviceStatusOffline {
|
||||
return latestUnavailable("手机已停用,请重新选择")
|
||||
}
|
||||
if service.OperatorUserID > 0 && !service.OperatorIsAdmin && (device.OwnerUserID == nil || *device.OwnerUserID != service.OperatorUserID) {
|
||||
return &ServiceError{Code: CodeDeviceOwnershipForbidden, Message: "该手机不属于当前采购员,无法使用其采集记录", Retryable: false}
|
||||
}
|
||||
var task models.CollectionTask
|
||||
query := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("device_id = ? AND source = ? AND status IN ?", r.DeviceID, models.CollectionTaskSourceAgentCurrentPage, []string{models.TaskStatusCompleted, models.TaskStatusCompletedPartial})
|
||||
if replacement != nil && !replacement.Preview {
|
||||
|
||||
@@ -63,6 +63,21 @@ func TestLatestCollectionOfflineAndOnlyIfUnlinked(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestLatestCollectionRejectsDeviceOwnedByAnotherUser(t *testing.T) {
|
||||
s, id, req, _ := latestFixture(t)
|
||||
owner := uint64(21)
|
||||
s.DB.Model(&models.AgentDevice{}).Where("id = ?", req.LatestCollection.DeviceID).Update("owner_user_id", owner)
|
||||
s.OperatorUserID = 22
|
||||
if _, err := s.LinkPDD(context.Background(), id, req); err == nil || errCode(t, err) != CodeDeviceOwnershipForbidden {
|
||||
t.Fatalf("expected ownership rejection: %v", err)
|
||||
}
|
||||
var product models.ShopeeProduct
|
||||
s.DB.First(&product, id)
|
||||
if product.PDDProductID != nil {
|
||||
t.Fatal("ownership rejection wrote association")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLatestCollectionRejectsInvalidContextWithoutWrite(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
|
||||
@@ -62,6 +62,8 @@ type ReplacementEligibility func(context.Context, *gorm.DB, uint64) error
|
||||
type Service struct {
|
||||
DB *gorm.DB
|
||||
ReplacementEligibility ReplacementEligibility
|
||||
OperatorUserID uint64
|
||||
OperatorIsAdmin bool
|
||||
}
|
||||
|
||||
func NewService(db *gorm.DB) *Service { return &Service{DB: db} }
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
package version_local
|
||||
|
||||
import (
|
||||
goautomigrations "go-admin/app/goauto/migrations"
|
||||
"go-admin/cmd/migrate/migration"
|
||||
common "go-admin/common/models"
|
||||
"gorm.io/gorm"
|
||||
"path/filepath"
|
||||
)
|
||||
|
||||
func init() {
|
||||
fileName := filepath.Base("1789800300000_device_owner.go")
|
||||
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateDeviceOwner)
|
||||
}
|
||||
func migrateDeviceOwner(tx *gorm.DB, version string) error {
|
||||
return tx.Transaction(func(db *gorm.DB) error {
|
||||
if err := goautomigrations.Migrate(db); err != nil {
|
||||
return err
|
||||
}
|
||||
return db.Create(&common.Migration{Version: version}).Error
|
||||
})
|
||||
}
|
||||
@@ -8,6 +8,14 @@ export function listDevices(params) {
|
||||
})
|
||||
}
|
||||
|
||||
export function listDeviceOwners() {
|
||||
return request({ url: '/api/admin/v1/devices/owners', method: 'get' })
|
||||
}
|
||||
|
||||
export function setDeviceOwner(deviceId, ownerUserId) {
|
||||
return request({ url: `/api/admin/v1/devices/${deviceId}/owner`, method: 'patch', data: { ownerUserId} })
|
||||
}
|
||||
|
||||
export function disableDevice(deviceId) {
|
||||
return request({
|
||||
url: `/api/admin/v1/devices/${deviceId}/disable`,
|
||||
|
||||
@@ -34,6 +34,9 @@
|
||||
<div class="device-model">{{ row.manufacturer }} {{ row.model }}</div>
|
||||
</template>
|
||||
</el-table-column>
|
||||
<el-table-column label="所属采购员" min-width="150">
|
||||
<template #default="{ row }">{{ ownerLabel(row) }}</template>
|
||||
</el-table-column>
|
||||
<el-table-column label="系统" min-width="130">
|
||||
<template #default="{ row }">Android {{ row.androidVersion }}</template>
|
||||
</el-table-column>
|
||||
@@ -70,6 +73,7 @@
|
||||
</el-table-column>
|
||||
<el-table-column v-if="isAdmin" label="操作" width="270" fixed="right">
|
||||
<template #default="{ row }">
|
||||
<el-button type="primary" link @click="openOwnerDialog(row)">归属</el-button>
|
||||
<el-button
|
||||
type="warning"
|
||||
link
|
||||
@@ -90,7 +94,7 @@
|
||||
>吊销 Token</el-button>
|
||||
</template>
|
||||
</el-table-column>
|
||||
</el-table>
|
||||
</el-table>
|
||||
|
||||
<pagination
|
||||
v-show="total > 0"
|
||||
@@ -100,6 +104,13 @@
|
||||
@pagination="getList"
|
||||
/>
|
||||
</el-card>
|
||||
<el-dialog v-model="ownerDialog.open" title="设置所属采购员" width="420px">
|
||||
<p v-if="ownerDialog.row">设备:{{ ownerDialog.row.name }}(#{{ ownerDialog.row.id }})</p>
|
||||
<el-select v-model="ownerDialog.ownerUserId" clearable placeholder="未归属" style="width: 100%" :loading="ownerDialog.loading">
|
||||
<el-option v-for="owner in ownerDialog.owners" :key="owner.userId" :label="owner.nickName || owner.username" :value="owner.userId" />
|
||||
</el-select>
|
||||
<template #footer><el-button @click="ownerDialog.open = false">取消</el-button><el-button type="primary" :loading="ownerDialog.saving" @click="saveOwner">保存</el-button></template>
|
||||
</el-dialog>
|
||||
<el-drawer v-model="releaseDrawer.open" title="Agent 版本管理" size="720px" destroy-on-close>
|
||||
<el-alert type="info" :closable="false" show-icon title="上传后需明确设为当前版本;APK 下载需要 Admin 或设备 Token,不提供公开静态地址。" />
|
||||
<el-form label-position="top" class="release-upload-form">
|
||||
@@ -122,7 +133,7 @@
|
||||
<script>
|
||||
import { ElMessage, ElMessageBox } from 'element-plus'
|
||||
import { Refresh, RefreshLeft, Search } from '@element-plus/icons-vue'
|
||||
import { disableDevice, listDevices, resetDeviceIdentity, revokeDeviceToken } from '@/api/goauto/devices'
|
||||
import { disableDevice, listDeviceOwners, listDevices, resetDeviceIdentity, revokeDeviceToken, setDeviceOwner } from '@/api/goauto/devices'
|
||||
import { downloadAgentAppRelease, listAgentAppReleases, setCurrentAgentAppRelease, uploadAgentAppRelease } from '@/api/goauto/agent-app-releases'
|
||||
import { createRequestId } from '@/utils/request-id'
|
||||
|
||||
@@ -137,7 +148,8 @@ export default {
|
||||
devices: [],
|
||||
total: 0,
|
||||
query: { page: 1, pageSize: 20, name: '', status: '' },
|
||||
releaseDrawer: { open: false, loading: false, uploading: false, progress: 0, file: null, notes: '', items: [] }
|
||||
releaseDrawer: { open: false, loading: false, uploading: false, progress: 0, file: null, notes: '', items: [] },
|
||||
ownerDialog: { open: false, loading: false, saving: false, row: null, ownerUserId: null, owners: [] }
|
||||
}
|
||||
},
|
||||
computed: {
|
||||
@@ -171,6 +183,9 @@ export default {
|
||||
statusType(status) {
|
||||
return { online: 'success', offline: 'info', disabled: 'danger' }[status] || 'info'
|
||||
},
|
||||
ownerLabel(row) { if (!row.ownerUserId) return '未归属'; const owner = this.ownerDialog.owners.find(item => item.userId === row.ownerUserId); return owner ? (owner.nickName || owner.username) : `用户 #${row.ownerUserId}` },
|
||||
async openOwnerDialog(row) { this.ownerDialog.row = row; this.ownerDialog.ownerUserId = row.ownerUserId || null; this.ownerDialog.open = true; this.ownerDialog.loading = true; try { this.ownerDialog.owners = (await listDeviceOwners()).data || [] } finally { this.ownerDialog.loading = false } },
|
||||
async saveOwner() { if (!this.ownerDialog.row) return; this.ownerDialog.saving = true; try { await setDeviceOwner(this.ownerDialog.row.id, this.ownerDialog.ownerUserId); ElMessage.success('设备归属已更新'); this.ownerDialog.open = false; await this.getList() } finally { this.ownerDialog.saving = false } },
|
||||
async openReleases() { this.releaseDrawer.open = true; await this.loadReleases() },
|
||||
async loadReleases() { this.releaseDrawer.loading = true; try { const response = await listAgentAppReleases({ page: 1, pageSize: 100 }); this.releaseDrawer.items = response.data.items } finally { this.releaseDrawer.loading = false } },
|
||||
onAPKChange(file) { this.releaseDrawer.file = file.raw }, onAPKRemove() { this.releaseDrawer.file = null },
|
||||
|
||||
Reference in New Issue
Block a user