From 696e0bca5273f4e8ca2bdabfe790909b2d2067ad Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Fri, 28 Aug 2026 17:33:15 +0800 Subject: [PATCH] feat(#129): add PDD product replacement model --- docs/02-architecture-and-code-map.md | 12 +- docs/03-business-rules-and-glossary.md | 13 +- server/app/admin/router/init_router.go | 2 + server/app/goauto/access/purchaser.go | 2 + server/app/goauto/access/purchaser_test.go | 18 + server/app/goauto/migrations/migrate.go | 2 + server/app/goauto/migrations/migrate_test.go | 26 ++ server/app/goauto/models/replacement.go | 126 +++++++ server/app/goauto/replacement/handler.go | 115 ++++++ server/app/goauto/replacement/handler_test.go | 58 +++ server/app/goauto/replacement/router.go | 15 + server/app/goauto/replacement/service.go | 335 ++++++++++++++++++ server/app/goauto/replacement/service_test.go | 281 +++++++++++++++ .../1787885000000_pdd_product_replacement.go | 25 ++ 14 files changed, 1026 insertions(+), 4 deletions(-) create mode 100644 server/app/goauto/models/replacement.go create mode 100644 server/app/goauto/replacement/handler.go create mode 100644 server/app/goauto/replacement/handler_test.go create mode 100644 server/app/goauto/replacement/router.go create mode 100644 server/app/goauto/replacement/service.go create mode 100644 server/app/goauto/replacement/service_test.go create mode 100644 server/cmd/migrate/migration/version-local/1787885000000_pdd_product_replacement.go diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index 35f620f..9be7041 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -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: f82f510e47cf3e74a3510a859ca17bddd2833362 -synchronized_at: 2026-08-28T06:51:01Z +wiki_revision: d585a906857a46752788c9f25d3ac0361739dada +synchronized_at: 2026-08-28T09:27:21Z # 架构与代码地图 @@ -226,3 +226,11 @@ Android Portal/Agent - `server/app/goauto/sybclient/` 提供 SYB 读取与三个入库码写接口;匹配器只依赖只读接口,不能取得写能力。真实写入只由用户确认后的回写批次触发。 - Admin API 注册在 `server/app/goauto/sybinnercode/router.go`,页面位于 `web/src/views/goauto/syb-inner-codes/`,菜单入口为“档口入库码”。 - 服务启动时,未开始的 `queued` 记录释放回可回写,已经进入远端动作的 `applying` 记录转为 `needs_check`,不自动重放写请求。 + +## PDD 商品替换数据模型(#129) + +- `server/app/goauto/models/replacement.go` 定义审计主表 `pdd_product_replacement` 与分项表 `pdd_product_replacement_item`;版本迁移为 `1787885000000_pdd_product_replacement.go`,只新增表,不修改既有表和数据。 +- 主表用 `origin_type + origin_task_id` 区分失败采集任务和失败采购任务来源,以 `target_collection_task_id` 保存替代商品采集证据,以 `created_by_device_id` 保存发起设备。当前设备模型没有操作人绑定,因此该字段只能追溯到设备,不能追溯到采购员账号。 +- 主表的可空 `active_slot` 与 `source_product_id` 组成唯一索引,使同一源商品跨 SQLite、MySQL 和 PostgreSQL 同时最多存在一条 `active` 替换记录;历史记录使用 `superseded`,不软删除、不提供删除接口。 +- 分项表按 `replacement_id + shopee_product_id` 唯一保存本次实际影响集合,并独立记录 `matching / matched / manual_required`、匹配来源、置信度和持久 worker 的尝试/错误信息。某条采购任务能否续做只能读取对应虾皮商品的分项状态,不能读取主表总体进度。 +- `server/app/goauto/replacement/` 提供内部登记、幂等冲突校验、来源/采集证据校验、目标有效性、活动记录唯一性和环检测,以及只读审计查询;管理端只读路由为 `GET /api/admin/v1/pdd-product-replacements` 与 `GET /api/admin/v1/pdd-product-replacements/{replacementId}`。Agent 没有直接写入该领域的 HTTP 权限。 diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index a380dca..56009f5 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -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: 0798ea888bbc55fd434e9ac31c539da7a4efea5a -synchronized_at: 2026-08-28T06:14:23Z +wiki_revision: b31ebf8888da7c3d9b71520938b104ff9e66511f +synchronized_at: 2026-08-28T09:27:33Z # 业务规则与术语 @@ -298,3 +298,12 @@ synchronized_at: 2026-08-28T06:14:23Z - 回写前弹窗展示业务记录数、入库码总数、预计占位明细数和替换旧码数。只有 `ready` 记录可提交;所有服务实例共用数据库租约全局串行执行。 - 每个远端写动作只发送一次,动作前重读整张货运单并校验匹配计划未漂移;每件写入后重读确认目标码唯一位于预期明细。超时、5xx、响应无法确认或服务在写入期间重启时转为 `needs_check`,禁止自动重试;“只读复核”只能读取远端状态。 - 列表支持勾选后批量物理删除。删除必须再次确认;选中项包含 `queued`、`applying` 或 `needs_check` 时整批拒绝,不做部分删除。其余选中业务记录、逐件码、计划和终态执行证据在同一事务中物理删除;已写入 SYB 的远端值不会撤销。 + +## PDD 失效或售罄商品替换(#129) + +- 商品替换是“失效的旧 PDD 商品 → 已成功采集的替代 PDD 商品”的即时生效审计事实,可由失败采集任务或失败采购任务发起;Agent 只提交明确的来源类型、任务编号和采集结果,服务端负责最终校验和登记。 +- 替换记录不删除。同一源商品同时只允许一条生效记录;替代商品必须为 `active`,不得与源商品相同,不得形成替换环,并必须由同一发起设备的一条 `completed` 或 `completed_partial` 采集任务证明。 +- 一个 PDD 商品可能被多个虾皮商品共用,因此规格匹配按受影响的每个虾皮商品分项记录。主表只表达总体进度;Agent 展示和“继续采购”资格必须读取当前任务对应虾皮商品的分项状态。 +- 选错替代商品时,连续执行 B→C 不等于撤销 A→B,因为 B 可能还关联其他虾皮商品。正确纠错语义是把原 A→B 记录置为 `superseded`,再建立 A→C,并只处理原记录分项中冻结的影响集合。 +- 创建请求按 `create_request_id` 幂等;重放时源商品、替代商品、来源类型、来源任务、采集证据和发起设备必须一致,否则返回幂等冲突,不能静默覆盖。 +- `created_by_device_id` 只代表 Agent 设备。当前系统没有设备到采购员账号的绑定,多人多机场景若需要个人责任追踪,必须另建工单实现设备绑定操作员。 diff --git a/server/app/admin/router/init_router.go b/server/app/admin/router/init_router.go index 45a5695..9a82299 100644 --- a/server/app/admin/router/init_router.go +++ b/server/app/admin/router/init_router.go @@ -10,6 +10,7 @@ import ( goautodevice "go-admin/app/goauto/device" goautoproduct "go-admin/app/goauto/product" goautopurchase "go-admin/app/goauto/purchase" + goautoreplacement "go-admin/app/goauto/replacement" goautorule "go-admin/app/goauto/rule" goautoshopeeproduct "go-admin/app/goauto/shopeeproduct" goautosybimport "go-admin/app/goauto/sybimport" @@ -54,6 +55,7 @@ func InitRouter() { goautotask.InitRouter(r, authMiddleware) goautopurchase.InitRouter(r, authMiddleware) goautoproduct.InitRouter(r, authMiddleware) + goautoreplacement.InitRouter(r, authMiddleware) goautorule.InitRouter(r, authMiddleware) goautoshopeeproduct.InitRouter(r, authMiddleware) goautosybimport.InitRouter(r, authMiddleware) diff --git a/server/app/goauto/access/purchaser.go b/server/app/goauto/access/purchaser.go index fa87b2a..24f44ea 100644 --- a/server/app/goauto/access/purchaser.go +++ b/server/app/goauto/access/purchaser.go @@ -23,6 +23,8 @@ var AdminAPIs = []APIPermission{ {"新增 PDD 商品", "/api/admin/v1/pdd-products", "POST", true}, {"查看 PDD 商品详情", "/api/admin/v1/pdd-products/:productId", "GET", true}, {"修改 PDD 商品", "/api/admin/v1/pdd-products/:productId", "PATCH", true}, + {"查看 PDD 商品替换记录", "/api/admin/v1/pdd-product-replacements", "GET", true}, + {"查看 PDD 商品替换详情", "/api/admin/v1/pdd-product-replacements/:replacementId", "GET", true}, {"查看虾皮商品", "/api/admin/v1/shopee-products", "GET", true}, {"新增虾皮商品", "/api/admin/v1/shopee-products", "POST", true}, diff --git a/server/app/goauto/access/purchaser_test.go b/server/app/goauto/access/purchaser_test.go index 12338d0..0b96088 100644 --- a/server/app/goauto/access/purchaser_test.go +++ b/server/app/goauto/access/purchaser_test.go @@ -26,3 +26,21 @@ func TestPurchaserExcludesAdministratorOperations(t *testing.T) { } } } + +func TestPurchaserMayOnlyReadReplacementAudit(t *testing.T) { + want := map[string]bool{ + "GET /api/admin/v1/pdd-product-replacements": false, + "GET /api/admin/v1/pdd-product-replacements/:replacementId": false, + } + for _, permission := range PurchaserAPIs() { + key := permission.Method + " " + permission.Path + if _, ok := want[key]; ok { + want[key] = true + } + } + for permission, found := range want { + if !found { + t.Fatalf("missing purchaser replacement audit permission: %s", permission) + } + } +} diff --git a/server/app/goauto/migrations/migrate.go b/server/app/goauto/migrations/migrate.go index 0191460..46ef634 100644 --- a/server/app/goauto/migrations/migrate.go +++ b/server/app/goauto/migrations/migrate.go @@ -41,6 +41,8 @@ func MigratedModels() []any { &models.CollectionColorPrice{}, &models.CollectionSKU{}, &models.CollectionSKUValue{}, + &models.PDDProductReplacement{}, + &models.PDDProductReplacementItem{}, } } diff --git a/server/app/goauto/migrations/migrate_test.go b/server/app/goauto/migrations/migrate_test.go index 0ffc41a..c7f728e 100644 --- a/server/app/goauto/migrations/migrate_test.go +++ b/server/app/goauto/migrations/migrate_test.go @@ -65,6 +65,7 @@ func TestMigrationIsIdempotentAndHasExpectedTables(t *testing.T) { for _, table := range []string{ "agent_device", "pdd_product", "shopee_product", "syb_product", "collection_rule", "agent_manual_collection_setting", "collection_task", "collection_dimension", "collection_dimension_value", "collection_sku", "collection_sku_value", + "pdd_product_replacement", "pdd_product_replacement_item", } { if !db.Migrator().HasTable(table) { t.Errorf("expected table %s", table) @@ -77,6 +78,31 @@ func TestMigrationIsIdempotentAndHasExpectedTables(t *testing.T) { } } +func TestReplacementMigrationPreservesExistingRows(t *testing.T) { + 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 := db.AutoMigrate(&models.PDDProduct{}); err != nil { + t.Fatal(err) + } + product := models.PDDProduct{GoodsID: "700000000099", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=700000000099", Status: "active"} + if err := db.Create(&product).Error; err != nil { + t.Fatal(err) + } + if err := migrations.Migrate(db); err != nil { + t.Fatalf("migrate existing database: %v", err) + } + var count int64 + if err := db.Model(&models.PDDProduct{}).Where("id = ? AND goods_id = ?", product.ID, product.GoodsID).Count(&count).Error; err != nil || count != 1 { + t.Fatalf("existing product count=%d err=%v", count, err) + } + if !db.Migrator().HasTable(&models.PDDProductReplacement{}) || !db.Migrator().HasTable(&models.PDDProductReplacementItem{}) { + t.Fatal("replacement tables missing after existing-database migration") + } +} + func TestMySQLCompositeIndexesFitInnoDBLimit(t *testing.T) { parsed, err := schema.Parse(&models.CollectionSKU{}, &sync.Map{}, schema.NamingStrategy{}) if err != nil { diff --git a/server/app/goauto/models/replacement.go b/server/app/goauto/models/replacement.go new file mode 100644 index 0000000..2b03334 --- /dev/null +++ b/server/app/goauto/models/replacement.go @@ -0,0 +1,126 @@ +package models + +import ( + "fmt" + "time" + + "gorm.io/gorm" +) + +const ( + ReplacementOriginCollection = "collection" + ReplacementOriginPurchase = "purchase" + + ReplacementStatusActive = "active" + ReplacementStatusSuperseded = "superseded" + + ReplacementMappingMatching = "matching" + ReplacementMappingCompleted = "completed" + ReplacementMappingCompletedPartial = "completed_partial" + + ReplacementItemMappingMatching = "matching" + ReplacementItemMappingMatched = "matched" + ReplacementItemMappingManualRequired = "manual_required" + + ReplacementItemSourceAIMatch = "ai_match" + ReplacementItemSourceExactMatch = "exact_match" + ReplacementItemSourceManualMapping = "manual_mapping" +) + +// PDDProductReplacement is an immutable audit record for one product +// replacement. ActiveSlot is a nullable uniqueness guard: exactly one active +// replacement may exist for a source product across SQLite, MySQL and +// PostgreSQL, while superseded history remains append-only. +type PDDProductReplacement struct { + ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"` + + SourceProductID uint64 `json:"sourceProductId" gorm:"not null;index;uniqueIndex:ux_pdd_product_replacement_active,priority:1"` + SourceProduct PDDProduct `json:"-" gorm:"foreignKey:SourceProductID;constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"` + TargetProductID uint64 `json:"targetProductId" gorm:"not null;index"` + TargetProduct PDDProduct `json:"-" gorm:"foreignKey:TargetProductID;constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"` + + OriginType string `json:"originType" gorm:"size:16;not null;index:ix_pdd_product_replacement_origin,priority:1;check:ck_pdd_product_replacement_origin_type,origin_type IN ('collection','purchase')"` + OriginTaskID uint64 `json:"originTaskId" gorm:"not null;index:ix_pdd_product_replacement_origin,priority:2"` + + TargetCollectionTaskID uint64 `json:"targetCollectionTaskId" gorm:"not null;index"` + TargetCollectionTask CollectionTask `json:"-" gorm:"foreignKey:TargetCollectionTaskID;constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"` + CreatedByDeviceID uint64 `json:"createdByDeviceId" gorm:"not null;index"` + CreatedByDevice AgentDevice `json:"-" gorm:"foreignKey:CreatedByDeviceID;constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"` + + Status string `json:"status" gorm:"size:16;not null;index;check:ck_pdd_product_replacement_status,status IN ('active','superseded')"` + ActiveSlot *uint8 `json:"-" gorm:"uniqueIndex:ux_pdd_product_replacement_active,priority:2;check:ck_pdd_product_replacement_active_slot,(status = 'active' AND active_slot = 1) OR (status = 'superseded' AND active_slot IS NULL)"` + MappingStatus string `json:"mappingStatus" gorm:"size:24;not null;index;check:ck_pdd_product_replacement_mapping_status,mapping_status IN ('matching','completed','completed_partial')"` + + CreateRequestID string `json:"-" gorm:"size:36;not null;uniqueIndex:ux_pdd_product_replacement_create_request"` + CreatedAt time.Time `json:"createdAt"` + MappingUpdatedAt time.Time `json:"mappingUpdatedAt" gorm:"not null"` +} + +func (PDDProductReplacement) TableName() string { return "pdd_product_replacement" } + +func (replacement *PDDProductReplacement) BeforeCreate(_ *gorm.DB) error { + if replacement.Status == "" { + replacement.Status = ReplacementStatusActive + } + if replacement.MappingStatus == "" { + replacement.MappingStatus = ReplacementMappingMatching + } + if replacement.MappingUpdatedAt.IsZero() { + replacement.MappingUpdatedAt = time.Now().UTC() + } + return replacement.syncActiveSlot() +} + +func (replacement *PDDProductReplacement) BeforeSave(_ *gorm.DB) error { + return replacement.syncActiveSlot() +} + +func (replacement *PDDProductReplacement) SetStatus(status string) error { + replacement.Status = status + return replacement.syncActiveSlot() +} + +func (replacement *PDDProductReplacement) syncActiveSlot() error { + one := uint8(1) + switch replacement.Status { + case ReplacementStatusActive: + replacement.ActiveSlot = &one + case ReplacementStatusSuperseded: + replacement.ActiveSlot = nil + default: + return fmt.Errorf("unsupported replacement status %q", replacement.Status) + } + return nil +} + +// PDDProductReplacementItem freezes the Shopee products affected by one +// replacement and tracks matching independently for every product. The parent +// aggregate status must never be used to decide whether one purchase may +// continue. +type PDDProductReplacementItem struct { + ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"` + + ReplacementID uint64 `json:"replacementId" gorm:"not null;index;uniqueIndex:ux_pdd_product_replacement_item,priority:1"` + Replacement PDDProductReplacement `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"` + ShopeeProductID uint64 `json:"shopeeProductId" gorm:"not null;index;uniqueIndex:ux_pdd_product_replacement_item,priority:2"` + ShopeeProduct ShopeeProduct `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:RESTRICT"` + + MappingStatus string `json:"mappingStatus" gorm:"size:24;not null;index;check:ck_pdd_product_replacement_item_status,mapping_status IN ('matching','matched','manual_required')"` + Source *string `json:"source,omitempty" gorm:"size:24;check:ck_pdd_product_replacement_item_source,source IS NULL OR source IN ('ai_match','exact_match','manual_mapping')"` + Confidence *float64 `json:"confidence,omitempty" gorm:"check:ck_pdd_product_replacement_item_confidence,confidence IS NULL OR (confidence >= 0 AND confidence <= 1)"` + Reason *string `json:"reason,omitempty" gorm:"size:500"` + + AttemptCount int `json:"attemptCount" gorm:"not null;default:0;check:ck_pdd_product_replacement_item_attempt_count,attempt_count >= 0"` + LastErrorCode *string `json:"lastErrorCode,omitempty" gorm:"size:64"` + LastErrorAt *time.Time `json:"lastErrorAt,omitempty"` + UpdatedAt time.Time `json:"updatedAt"` +} + +func (PDDProductReplacementItem) TableName() string { return "pdd_product_replacement_item" } + +func (item *PDDProductReplacementItem) BeforeCreate(_ *gorm.DB) error { + if item.MappingStatus == "" { + item.MappingStatus = ReplacementItemMappingMatching + } + return nil +} diff --git a/server/app/goauto/replacement/handler.go b/server/app/goauto/replacement/handler.go new file mode 100644 index 0000000..f17de5a --- /dev/null +++ b/server/app/goauto/replacement/handler.go @@ -0,0 +1,115 @@ +package replacement + +import ( + "errors" + "net/http" + "strconv" + "strings" + + "github.com/gin-gonic/gin" + "github.com/go-admin-team/go-admin-core/sdk/pkg" + "gorm.io/gorm" +) + +type Handler struct{ DB *gorm.DB } + +func (handler Handler) List(context *gin.Context) { + page, err := positiveQuery(context.Query("page"), 1) + if err != nil { + writeError(context, fail(CodeInvalidRequest, "page 无效")) + return + } + pageSize, err := positiveQuery(context.Query("pageSize"), 20) + if err != nil || pageSize > 100 { + writeError(context, fail(CodeInvalidRequest, "pageSize 无效")) + return + } + sourceProductID, err := optionalUint(context.Query("sourceProductId")) + if err != nil { + writeError(context, fail(CodeInvalidRequest, "sourceProductId 无效")) + return + } + service, ok := handler.service(context) + if !ok { + return + } + result, err := service.List(context.Request.Context(), ListRequest{Page: page, PageSize: pageSize, SourceProductID: sourceProductID, Status: strings.TrimSpace(context.Query("status"))}) + if err != nil { + writeError(context, err) + return + } + context.Header("Cache-Control", "no-store") + context.JSON(http.StatusOK, gin.H{"data": result}) +} + +func (handler Handler) Detail(context *gin.Context) { + id, err := strconv.ParseUint(context.Param("replacementId"), 10, 64) + if err != nil || id == 0 { + writeError(context, fail(CodeInvalidRequest, "replacementId 无效")) + return + } + service, ok := handler.service(context) + if !ok { + return + } + result, err := service.Get(context.Request.Context(), id) + if err != nil { + writeError(context, err) + return + } + context.Header("Cache-Control", "no-store") + context.JSON(http.StatusOK, gin.H{"data": result}) +} + +func (handler Handler) service(context *gin.Context) (*Service, bool) { + db := handler.DB + if db == nil { + var err error + db, err = pkg.GetOrm(context) + if err != nil { + writeError(context, internal(err)) + return nil, false + } + } + return NewService(db), true +} + +func writeError(context *gin.Context, err error) { + status := http.StatusInternalServerError + code, message, retryable := CodeInternal, "替换记录处理失败", true + var serviceError *ServiceError + if errors.As(err, &serviceError) { + code, message, retryable = serviceError.Code, serviceError.Message, serviceError.Retryable + } + switch code { + case CodeInvalidRequest: + status = http.StatusUnprocessableEntity + case CodeNotFound, CodeProductNotFound, CodeOriginTaskNotFound: + status = http.StatusNotFound + case CodeTargetNotActive, CodeSameProduct, CodeOriginTaskInvalid, CodeCollectionProofInvalid, CodeSourceAlreadyReplaced, CodeCycleDetected, CodeIdempotencyConflict: + status = http.StatusConflict + } + context.JSON(status, gin.H{"code": code, "message": message, "retryable": retryable}) +} + +func positiveQuery(raw string, fallback int) (int, error) { + if strings.TrimSpace(raw) == "" { + return fallback, nil + } + value, err := strconv.Atoi(raw) + if err != nil || value < 1 { + return 0, errors.New("invalid positive number") + } + return value, nil +} + +func optionalUint(raw string) (uint64, error) { + if strings.TrimSpace(raw) == "" { + return 0, nil + } + value, err := strconv.ParseUint(raw, 10, 64) + if err != nil || value == 0 { + return 0, errors.New("invalid uint") + } + return value, nil +} diff --git a/server/app/goauto/replacement/handler_test.go b/server/app/goauto/replacement/handler_test.go new file mode 100644 index 0000000..1388b59 --- /dev/null +++ b/server/app/goauto/replacement/handler_test.go @@ -0,0 +1,58 @@ +package replacement + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + + "github.com/gin-gonic/gin" +) + +func TestReadOnlyHandlersListAndReturnItems(t *testing.T) { + gin.SetMode(gin.TestMode) + fixture := seedReplacementFixture(t) + created, err := fixture.service.Register(t.Context(), fixture.request()) + if err != nil { + t.Fatal(err) + } + engine := gin.New() + handler := Handler{DB: fixture.db} + engine.GET("/replacements", handler.List) + engine.GET("/replacements/:replacementId", handler.Detail) + + listRecorder := httptest.NewRecorder() + engine.ServeHTTP(listRecorder, httptest.NewRequest(http.MethodGet, "/replacements?page=1&pageSize=20&status=active", nil)) + if listRecorder.Code != http.StatusOK || listRecorder.Header().Get("Cache-Control") != "no-store" { + t.Fatalf("list status=%d body=%s", listRecorder.Code, listRecorder.Body.String()) + } + var listBody struct { + Data ListResponse `json:"data"` + } + if err := json.Unmarshal(listRecorder.Body.Bytes(), &listBody); err != nil || listBody.Data.Total != 1 { + t.Fatalf("list body=%s err=%v", listRecorder.Body.String(), err) + } + + detailRecorder := httptest.NewRecorder() + engine.ServeHTTP(detailRecorder, httptest.NewRequest(http.MethodGet, "/replacements/"+jsonNumber(created.Replacement.ID), nil)) + if detailRecorder.Code != http.StatusOK || !json.Valid(detailRecorder.Body.Bytes()) { + t.Fatalf("detail status=%d body=%s", detailRecorder.Code, detailRecorder.Body.String()) + } +} + +func TestReadOnlyHandlersRejectInvalidFilters(t *testing.T) { + gin.SetMode(gin.TestMode) + fixture := seedReplacementFixture(t) + engine := gin.New() + engine.GET("/replacements", Handler{DB: fixture.db}.List) + recorder := httptest.NewRecorder() + engine.ServeHTTP(recorder, httptest.NewRequest(http.MethodGet, "/replacements?pageSize=101", nil)) + if recorder.Code != http.StatusUnprocessableEntity { + t.Fatalf("status=%d body=%s", recorder.Code, recorder.Body.String()) + } +} + +func jsonNumber(value uint64) string { + encoded, _ := json.Marshal(value) + return string(encoded) +} diff --git a/server/app/goauto/replacement/router.go b/server/app/goauto/replacement/router.go new file mode 100644 index 0000000..f06d5f6 --- /dev/null +++ b/server/app/goauto/replacement/router.go @@ -0,0 +1,15 @@ +package replacement + +import ( + "go-admin/common/middleware" + + "github.com/gin-gonic/gin" + jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth" +) + +func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) { + handler := Handler{} + admin := engine.Group("/api/admin/v1/pdd-product-replacements").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()) + admin.GET("", handler.List) + admin.GET("/:replacementId", handler.Detail) +} diff --git a/server/app/goauto/replacement/service.go b/server/app/goauto/replacement/service.go new file mode 100644 index 0000000..604fe11 --- /dev/null +++ b/server/app/goauto/replacement/service.go @@ -0,0 +1,335 @@ +package replacement + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "go-admin/app/goauto/models" + + "github.com/google/uuid" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +const ( + CodeInvalidRequest = "REPLACEMENT_INVALID_REQUEST" + CodeProductNotFound = "REPLACEMENT_PRODUCT_NOT_FOUND" + CodeTargetNotActive = "REPLACEMENT_TARGET_NOT_ACTIVE" + CodeSameProduct = "REPLACEMENT_SAME_PRODUCT" + CodeOriginTaskNotFound = "REPLACEMENT_ORIGIN_TASK_NOT_FOUND" + CodeOriginTaskInvalid = "REPLACEMENT_ORIGIN_TASK_INVALID" + CodeCollectionProofInvalid = "REPLACEMENT_COLLECTION_PROOF_INVALID" + CodeSourceAlreadyReplaced = "REPLACEMENT_SOURCE_ALREADY_REPLACED" + CodeCycleDetected = "REPLACEMENT_CYCLE_DETECTED" + CodeIdempotencyConflict = "REPLACEMENT_IDEMPOTENCY_CONFLICT" + CodeNotFound = "REPLACEMENT_NOT_FOUND" + CodeInternal = "REPLACEMENT_INTERNAL" +) + +type ServiceError struct { + Code string + Message string + Retryable bool + Cause error +} + +func (e *ServiceError) Error() string { + if e.Cause != nil { + return fmt.Sprintf("%s: %v", e.Message, e.Cause) + } + return e.Message +} + +func (e *ServiceError) Unwrap() error { return e.Cause } + +type RegisterRequest struct { + RequestID string + SourceProductID uint64 + TargetProductID uint64 + OriginType string + OriginTaskID uint64 + TargetCollectionTaskID uint64 + CreatedByDeviceID uint64 +} + +type Record struct { + Replacement models.PDDProductReplacement `json:"replacement"` + Items []models.PDDProductReplacementItem `json:"items"` + Replayed bool `json:"replayed"` +} + +type ListRequest struct { + Page int + PageSize int + SourceProductID uint64 + Status string +} + +type ListResponse struct { + Items []models.PDDProductReplacement `json:"items"` + Total int64 `json:"total"` + Page int `json:"page"` + PageSize int `json:"pageSize"` +} + +type Service struct { + DB *gorm.DB + Now func() time.Time +} + +func NewService(db *gorm.DB) *Service { + return &Service{DB: db, Now: func() time.Time { return time.Now().UTC() }} +} + +// Register writes the audit fact only. #131 owns the activation transaction +// and item creation; callers must not interpret this method as changing Shopee +// links or purchase tasks. +func (service *Service) Register(ctx context.Context, request RegisterRequest) (Record, error) { + if err := validateRegisterRequest(request); err != nil { + return Record{}, err + } + if service.DB == nil { + return Record{}, internal(errors.New("database is nil")) + } + + var record models.PDDProductReplacement + replayed := false + err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var existing models.PDDProductReplacement + if err := tx.Where("create_request_id = ?", request.RequestID).First(&existing).Error; err == nil { + if !sameRequest(existing, request) { + return fail(CodeIdempotencyConflict, "requestId 已用于其他替换请求") + } + record, replayed = existing, true + return nil + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return internal(err) + } + + if err := validateProducts(tx, request); err != nil { + return err + } + if err := validateOrigin(tx, request); err != nil { + return err + } + if err := validateCollectionProof(tx, request); err != nil { + return err + } + if err := ensureNoActiveReplacement(tx, request.SourceProductID); err != nil { + return err + } + if err := ensureNoCycle(tx, request.SourceProductID, request.TargetProductID); err != nil { + return err + } + + now := service.Now() + record = models.PDDProductReplacement{ + SourceProductID: request.SourceProductID, TargetProductID: request.TargetProductID, + OriginType: request.OriginType, OriginTaskID: request.OriginTaskID, + TargetCollectionTaskID: request.TargetCollectionTaskID, CreatedByDeviceID: request.CreatedByDeviceID, + Status: models.ReplacementStatusActive, MappingStatus: models.ReplacementMappingMatching, + CreateRequestID: request.RequestID, CreatedAt: now, MappingUpdatedAt: now, + } + if err := tx.Create(&record).Error; err != nil { + if isUniqueConstraint(err) { + return fail(CodeSourceAlreadyReplaced, "该失效商品已有生效中的替换记录") + } + return internal(err) + } + return nil + }) + if err != nil { + return Record{}, err + } + return Record{Replacement: record, Items: []models.PDDProductReplacementItem{}, Replayed: replayed}, nil +} + +func (service *Service) Get(ctx context.Context, replacementID uint64) (Record, error) { + if replacementID == 0 || service.DB == nil { + return Record{}, fail(CodeInvalidRequest, "replacementId 无效") + } + var record models.PDDProductReplacement + if err := service.DB.WithContext(ctx).First(&record, replacementID).Error; errors.Is(err, gorm.ErrRecordNotFound) { + return Record{}, fail(CodeNotFound, "替换记录不存在") + } else if err != nil { + return Record{}, internal(err) + } + var items []models.PDDProductReplacementItem + if err := service.DB.WithContext(ctx).Where("replacement_id = ?", record.ID).Order("id ASC").Find(&items).Error; err != nil { + return Record{}, internal(err) + } + return Record{Replacement: record, Items: items}, nil +} + +func (service *Service) ActiveBySource(ctx context.Context, sourceProductID uint64) (Record, error) { + if sourceProductID == 0 || service.DB == nil { + return Record{}, fail(CodeInvalidRequest, "sourceProductId 无效") + } + var record models.PDDProductReplacement + if err := service.DB.WithContext(ctx).Where("source_product_id = ? AND status = ?", sourceProductID, models.ReplacementStatusActive).First(&record).Error; errors.Is(err, gorm.ErrRecordNotFound) { + return Record{}, fail(CodeNotFound, "替换记录不存在") + } else if err != nil { + return Record{}, internal(err) + } + return service.Get(ctx, record.ID) +} + +func (service *Service) List(ctx context.Context, request ListRequest) (ListResponse, error) { + if service.DB == nil { + return ListResponse{}, internal(errors.New("database is nil")) + } + if request.Page < 1 { + request.Page = 1 + } + if request.PageSize < 1 { + request.PageSize = 20 + } + if request.PageSize > 100 { + request.PageSize = 100 + } + if request.Status != "" && request.Status != models.ReplacementStatusActive && request.Status != models.ReplacementStatusSuperseded { + return ListResponse{}, fail(CodeInvalidRequest, "status 无效") + } + query := service.DB.WithContext(ctx).Model(&models.PDDProductReplacement{}) + if request.SourceProductID != 0 { + query = query.Where("source_product_id = ?", request.SourceProductID) + } + if request.Status != "" { + query = query.Where("status = ?", request.Status) + } + var total int64 + if err := query.Count(&total).Error; err != nil { + return ListResponse{}, internal(err) + } + items := make([]models.PDDProductReplacement, 0) + if err := query.Order("created_at DESC, id DESC").Offset((request.Page - 1) * request.PageSize).Limit(request.PageSize).Find(&items).Error; err != nil { + return ListResponse{}, internal(err) + } + return ListResponse{Items: items, Total: total, Page: request.Page, PageSize: request.PageSize}, nil +} + +func validateRegisterRequest(request RegisterRequest) error { + if _, err := uuid.Parse(strings.TrimSpace(request.RequestID)); err != nil { + return fail(CodeInvalidRequest, "requestId 无效") + } + if request.SourceProductID == 0 || request.TargetProductID == 0 || request.OriginTaskID == 0 || request.TargetCollectionTaskID == 0 || request.CreatedByDeviceID == 0 { + return fail(CodeInvalidRequest, "替换请求缺少必要字段") + } + if request.OriginType != models.ReplacementOriginCollection && request.OriginType != models.ReplacementOriginPurchase { + return fail(CodeInvalidRequest, "originType 无效") + } + if request.SourceProductID == request.TargetProductID { + return fail(CodeSameProduct, "替代商品不能与失效商品相同") + } + return nil +} + +func validateProducts(tx *gorm.DB, request RegisterRequest) error { + var products []models.PDDProduct + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id IN ?", []uint64{request.SourceProductID, request.TargetProductID}).Find(&products).Error; err != nil { + return internal(err) + } + if len(products) != 2 { + return fail(CodeProductNotFound, "失效商品或替代商品不存在") + } + for _, product := range products { + if product.ID == request.TargetProductID && product.Status != "active" { + return fail(CodeTargetNotActive, "替代商品不是可用状态") + } + } + return nil +} + +func validateOrigin(tx *gorm.DB, request RegisterRequest) error { + switch request.OriginType { + case models.ReplacementOriginCollection: + var task models.CollectionTask + if err := tx.First(&task, request.OriginTaskID).Error; errors.Is(err, gorm.ErrRecordNotFound) { + return fail(CodeOriginTaskNotFound, "来源采集任务不存在") + } else if err != nil { + return internal(err) + } + if task.Status != models.TaskStatusFailed || task.DeviceID == nil || *task.DeviceID != request.CreatedByDeviceID || task.PDDProductID == nil || *task.PDDProductID != request.SourceProductID { + return fail(CodeOriginTaskInvalid, "来源采集任务不满足替换条件") + } + case models.ReplacementOriginPurchase: + var task models.PurchaseTask + if err := tx.First(&task, request.OriginTaskID).Error; errors.Is(err, gorm.ErrRecordNotFound) { + return fail(CodeOriginTaskNotFound, "来源采购任务不存在") + } else if err != nil { + return internal(err) + } + if task.Status != models.PurchaseTaskStatusFailed || task.DeviceID == nil || *task.DeviceID != request.CreatedByDeviceID || task.PDDProductID != request.SourceProductID { + return fail(CodeOriginTaskInvalid, "来源采购任务不满足替换条件") + } + } + return nil +} + +func validateCollectionProof(tx *gorm.DB, request RegisterRequest) error { + var task models.CollectionTask + if err := tx.First(&task, request.TargetCollectionTaskID).Error; errors.Is(err, gorm.ErrRecordNotFound) { + return fail(CodeCollectionProofInvalid, "替代商品采集证据不存在") + } else if err != nil { + return internal(err) + } + completed := task.Status == models.TaskStatusCompleted || task.Status == models.TaskStatusCompletedPartial + if !completed || task.DeviceID == nil || *task.DeviceID != request.CreatedByDeviceID || task.PDDProductID == nil || *task.PDDProductID != request.TargetProductID { + return fail(CodeCollectionProofInvalid, "替代商品采集证据无效") + } + return nil +} + +func ensureNoActiveReplacement(tx *gorm.DB, sourceProductID uint64) error { + var count int64 + if err := tx.Model(&models.PDDProductReplacement{}).Where("source_product_id = ? AND status = ?", sourceProductID, models.ReplacementStatusActive).Count(&count).Error; err != nil { + return internal(err) + } + if count != 0 { + return fail(CodeSourceAlreadyReplaced, "该失效商品已有生效中的替换记录") + } + return nil +} + +func ensureNoCycle(tx *gorm.DB, sourceProductID, targetProductID uint64) error { + cursor := targetProductID + seen := map[uint64]bool{sourceProductID: true} + for hops := 0; hops < 256; hops++ { + if seen[cursor] { + return fail(CodeCycleDetected, "替换关系会形成循环") + } + seen[cursor] = true + var next models.PDDProductReplacement + err := tx.Where("source_product_id = ? AND status = ?", cursor, models.ReplacementStatusActive).First(&next).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil + } + if err != nil { + return internal(err) + } + cursor = next.TargetProductID + } + return fail(CodeCycleDetected, "替换关系过长或存在循环") +} + +func sameRequest(record models.PDDProductReplacement, request RegisterRequest) bool { + return record.SourceProductID == request.SourceProductID && record.TargetProductID == request.TargetProductID && + record.OriginType == request.OriginType && record.OriginTaskID == request.OriginTaskID && + record.TargetCollectionTaskID == request.TargetCollectionTaskID && record.CreatedByDeviceID == request.CreatedByDeviceID +} + +func isUniqueConstraint(err error) bool { + message := strings.ToLower(err.Error()) + return strings.Contains(message, "unique constraint") || strings.Contains(message, "duplicate entry") || strings.Contains(message, "duplicate key") +} + +func fail(code, message string) error { + return &ServiceError{Code: code, Message: message, Retryable: false} +} + +func internal(err error) error { + return &ServiceError{Code: CodeInternal, Message: "替换记录处理失败", Retryable: true, Cause: err} +} diff --git a/server/app/goauto/replacement/service_test.go b/server/app/goauto/replacement/service_test.go new file mode 100644 index 0000000..83287c2 --- /dev/null +++ b/server/app/goauto/replacement/service_test.go @@ -0,0 +1,281 @@ +package replacement + +import ( + "context" + "errors" + "fmt" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + "go-admin/app/goauto/migrations" + "go-admin/app/goauto/models" + + "github.com/google/uuid" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + "gorm.io/gorm/logger" +) + +type replacementFixture struct { + db *gorm.DB + service *Service + device models.AgentDevice + source models.PDDProduct + target models.PDDProduct + rule models.CollectionRule + origin models.CollectionTask + proof models.CollectionTask +} + +func replacementDB(t *testing.T) *gorm.DB { + t.Helper() + databasePath := filepath.ToSlash(filepath.Join(t.TempDir(), fmt.Sprintf("replacement-%s.db", uuid.NewString()))) + dsn := fmt.Sprintf("file:%s?_foreign_keys=on&_journal_mode=WAL&_busy_timeout=5000", databasePath) + db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)}) + if err != nil { + t.Fatal(err) + } + if sqlDB, err := db.DB(); err == nil { + sqlDB.SetMaxOpenConns(8) + t.Cleanup(func() { _ = sqlDB.Close() }) + } + if err := migrations.Migrate(db); err != nil { + t.Fatal(err) + } + return db +} + +func seedReplacementFixture(t *testing.T) replacementFixture { + t.Helper() + db := replacementDB(t) + now := time.Date(2026, 8, 28, 9, 0, 0, 0, time.UTC) + device := models.AgentDevice{ + InstallID: uuid.NewString(), Name: "PKG110", Manufacturer: "Google", Model: "Pixel", + AndroidVersion: "15", AgentVersion: "1", PDDVersion: "7", CapabilitiesJSON: "[]", + Status: models.DeviceStatusOnline, TokenDigest: strings.ReplaceAll(uuid.NewString(), "-", "") + strings.ReplaceAll(uuid.NewString(), "-", ""), TokenIssuedAt: now, + } + source := models.PDDProduct{GoodsID: "700000000001", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=700000000001", Status: "disabled", SpecsJSON: "[]"} + target := models.PDDProduct{GoodsID: "700000000002", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=700000000002", Status: "active", SpecsJSON: "[]"} + rule := models.CollectionRule{Name: "PDD", ContentJSON: `{}`} + for _, value := range []any{&device, &source, &target, &rule} { + if err := db.Create(value).Error; err != nil { + t.Fatalf("seed %T: %v", value, err) + } + } + origin := models.CollectionTask{ + PDDProductID: &source.ID, RuleID: rule.ID, DeviceID: &device.ID, Source: models.CollectionTaskSourceAgentCurrentPage, + Status: models.TaskStatusFailed, URLSnapshot: source.URL, GoodsIDSnapshot: source.GoodsID, RuleSnapshot: `{}`, + } + proof := models.CollectionTask{ + PDDProductID: &target.ID, RuleID: rule.ID, DeviceID: &device.ID, Source: models.CollectionTaskSourceAgentCurrentPage, + Status: models.TaskStatusCompleted, URLSnapshot: target.URL, GoodsIDSnapshot: target.GoodsID, RuleSnapshot: `{}`, + } + for _, value := range []any{&origin, &proof} { + if err := db.Create(value).Error; err != nil { + t.Fatalf("seed collection task: %v", err) + } + } + service := NewService(db) + service.Now = func() time.Time { return now } + return replacementFixture{db: db, service: service, device: device, source: source, target: target, rule: rule, origin: origin, proof: proof} +} + +func (fixture replacementFixture) request() RegisterRequest { + return RegisterRequest{ + RequestID: uuid.NewString(), SourceProductID: fixture.source.ID, TargetProductID: fixture.target.ID, + OriginType: models.ReplacementOriginCollection, OriginTaskID: fixture.origin.ID, + TargetCollectionTaskID: fixture.proof.ID, CreatedByDeviceID: fixture.device.ID, + } +} + +func replacementCode(err error) string { + var target *ServiceError + if errors.As(err, &target) { + return target.Code + } + return "" +} + +func TestRegisterAndQueryCollectionOrigin(t *testing.T) { + fixture := seedReplacementFixture(t) + request := fixture.request() + created, err := fixture.service.Register(context.Background(), request) + if err != nil { + t.Fatalf("register: %v", err) + } + if created.Replayed || created.Replacement.Status != models.ReplacementStatusActive || created.Replacement.MappingStatus != models.ReplacementMappingMatching { + t.Fatalf("unexpected record: %+v", created) + } + if created.Replacement.ActiveSlot == nil || *created.Replacement.ActiveSlot != 1 { + t.Fatal("active uniqueness guard missing") + } + queried, err := fixture.service.ActiveBySource(context.Background(), fixture.source.ID) + if err != nil || queried.Replacement.ID != created.Replacement.ID || queried.Items == nil { + t.Fatalf("query active: %+v, %v", queried, err) + } + + replay, err := fixture.service.Register(context.Background(), request) + if err != nil || !replay.Replayed || replay.Replacement.ID != created.Replacement.ID { + t.Fatalf("idempotent replay: %+v, %v", replay, err) + } + conflict := request + conflict.TargetProductID = fixture.source.ID + if replacementCode(mustRegisterError(fixture.service, conflict)) != CodeSameProduct { + t.Fatal("request validation must run before idempotency lookup for an invalid same-product request") + } + conflict = request + other := models.PDDProduct{GoodsID: "700000000003", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=700000000003", Status: "active", SpecsJSON: "[]"} + if err := fixture.db.Create(&other).Error; err != nil { + t.Fatal(err) + } + conflict.TargetProductID = other.ID + if replacementCode(mustRegisterError(fixture.service, conflict)) != CodeIdempotencyConflict { + t.Fatal("same requestId with different target must conflict") + } +} + +func TestRegisterSupportsPurchaseOrigin(t *testing.T) { + fixture := seedReplacementFixture(t) + shopee := models.ShopeeProduct{ShopeeItemID: "SP-1", PDDProductID: &fixture.source.ID, SpecsJSON: "[]", Currency: "CNY"} + if err := fixture.db.Create(&shopee).Error; err != nil { + t.Fatal(err) + } + syb := models.SYBProduct{OrderCode: "SYB-1", DetailID: 1, StockID: 1, ShopeeItemID: shopee.ShopeeItemID, ShopeeProductID: &shopee.ID, Quantity: 1, UnitPriceCent: 100, ParseStatus: models.SYBParseStatusSuccess, RawJSON: `{}`} + if err := fixture.db.Create(&syb).Error; err != nil { + t.Fatal(err) + } + purchase := models.PurchaseTask{ + SYBProductID: &syb.ID, ShopeeProductID: &shopee.ID, PDDProductID: fixture.source.ID, DeviceID: &fixture.device.ID, + ExecutionMode: models.PurchaseExecutionModeLive, Status: models.PurchaseTaskStatusFailed, + ShopeeItemIDSnapshot: shopee.ShopeeItemID, PDDURLSnapshot: fixture.source.URL, PDDGoodsIDSnapshot: fixture.source.GoodsID, + Quantity: 1, Currency: "CNY", RuleType: "pddPurchase", RuleSchemaVersion: 1, RuleSnapshot: `{}`, + CreateRequestID: uuid.NewString(), PaymentReviewStatus: models.PurchasePaymentReviewPending, + LogisticsStatus: models.PurchaseLogisticsStatusPending, WritebackStatus: models.PurchaseWritebackStatusNotSelected, + } + if err := fixture.db.Create(&purchase).Error; err != nil { + t.Fatal(err) + } + request := fixture.request() + request.OriginType, request.OriginTaskID = models.ReplacementOriginPurchase, purchase.ID + created, err := fixture.service.Register(context.Background(), request) + if err != nil || created.Replacement.OriginType != models.ReplacementOriginPurchase { + t.Fatalf("purchase origin: %+v, %v", created, err) + } +} + +func TestRegisterRejectsInvalidInputsAndCycle(t *testing.T) { + tests := []struct { + name string + mutate func(*replacementFixture, *RegisterRequest) + code string + }{ + {"same product", func(f *replacementFixture, r *RegisterRequest) { r.TargetProductID = r.SourceProductID }, CodeSameProduct}, + {"inactive target", func(f *replacementFixture, r *RegisterRequest) { + f.db.Model(&models.PDDProduct{}).Where("id = ?", f.target.ID).Update("status", "pending") + }, CodeTargetNotActive}, + {"missing proof", func(f *replacementFixture, r *RegisterRequest) { r.TargetCollectionTaskID = 999999 }, CodeCollectionProofInvalid}, + {"wrong origin device", func(f *replacementFixture, r *RegisterRequest) { r.CreatedByDeviceID++ }, CodeOriginTaskInvalid}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + fixture := seedReplacementFixture(t) + request := fixture.request() + test.mutate(&fixture, &request) + if code := replacementCode(mustRegisterError(fixture.service, request)); code != test.code { + t.Fatalf("code=%s want=%s", code, test.code) + } + }) + } + + fixture := seedReplacementFixture(t) + back := models.PDDProductReplacement{ + SourceProductID: fixture.target.ID, TargetProductID: fixture.source.ID, + OriginType: models.ReplacementOriginCollection, OriginTaskID: fixture.origin.ID, + TargetCollectionTaskID: fixture.proof.ID, CreatedByDeviceID: fixture.device.ID, + Status: models.ReplacementStatusActive, MappingStatus: models.ReplacementMappingMatching, + CreateRequestID: uuid.NewString(), MappingUpdatedAt: time.Now().UTC(), + } + if err := fixture.db.Create(&back).Error; err != nil { + t.Fatal(err) + } + if code := replacementCode(mustRegisterError(fixture.service, fixture.request())); code != CodeCycleDetected { + t.Fatalf("cycle code=%s", code) + } +} + +func TestOnlyOneActiveReplacementPerSourceUnderConcurrency(t *testing.T) { + fixture := seedReplacementFixture(t) + requests := []RegisterRequest{fixture.request(), fixture.request()} + start := make(chan struct{}) + errs := make(chan error, len(requests)) + var wait sync.WaitGroup + for _, request := range requests { + wait.Add(1) + go func(request RegisterRequest) { + defer wait.Done() + <-start + _, err := fixture.service.Register(context.Background(), request) + errs <- err + }(request) + } + close(start) + wait.Wait() + close(errs) + successes := 0 + var failures []string + for err := range errs { + if err == nil { + successes++ + } else { + failures = append(failures, err.Error()) + } + } + if successes != 1 { + t.Fatalf("successes=%d want=1 failures=%v", successes, failures) + } + var count int64 + if err := fixture.db.Model(&models.PDDProductReplacement{}).Where("source_product_id = ? AND status = ?", fixture.source.ID, models.ReplacementStatusActive).Count(&count).Error; err != nil || count != 1 { + t.Fatalf("active count=%d err=%v", count, err) + } +} + +func TestReplacementItemsKeepIndependentStatuses(t *testing.T) { + fixture := seedReplacementFixture(t) + created, err := fixture.service.Register(context.Background(), fixture.request()) + if err != nil { + t.Fatal(err) + } + products := []models.ShopeeProduct{ + {ShopeeItemID: "SP-A", PDDProductID: &fixture.source.ID, SpecsJSON: "[]", Currency: "CNY"}, + {ShopeeItemID: "SP-B", PDDProductID: &fixture.source.ID, SpecsJSON: "[]", Currency: "CNY"}, + } + for index := range products { + if err := fixture.db.Create(&products[index]).Error; err != nil { + t.Fatal(err) + } + } + source := models.ReplacementItemSourceExactMatch + items := []models.PDDProductReplacementItem{ + {ReplacementID: created.Replacement.ID, ShopeeProductID: products[0].ID, MappingStatus: models.ReplacementItemMappingMatched, Source: &source}, + {ReplacementID: created.Replacement.ID, ShopeeProductID: products[1].ID, MappingStatus: models.ReplacementItemMappingManualRequired}, + } + if err := fixture.db.Create(&items).Error; err != nil { + t.Fatal(err) + } + queried, err := fixture.service.Get(context.Background(), created.Replacement.ID) + if err != nil || len(queried.Items) != 2 || queried.Items[0].MappingStatus == queried.Items[1].MappingStatus { + t.Fatalf("items=%+v err=%v", queried.Items, err) + } + duplicate := models.PDDProductReplacementItem{ReplacementID: created.Replacement.ID, ShopeeProductID: products[0].ID} + if err := fixture.db.Create(&duplicate).Error; err == nil { + t.Fatal("duplicate replacement item was accepted") + } +} + +func mustRegisterError(service *Service, request RegisterRequest) error { + _, err := service.Register(context.Background(), request) + return err +} diff --git a/server/cmd/migrate/migration/version-local/1787885000000_pdd_product_replacement.go b/server/cmd/migrate/migration/version-local/1787885000000_pdd_product_replacement.go new file mode 100644 index 0000000..56de43d --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1787885000000_pdd_product_replacement.go @@ -0,0 +1,25 @@ +package version_local + +import ( + "runtime" + + goautomigrations "go-admin/app/goauto/migrations" + "go-admin/cmd/migrate/migration" + common "go-admin/common/models" + + "gorm.io/gorm" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migratePDDProductReplacement) +} + +func migratePDDProductReplacement(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 + }) +}