feat(#129): add PDD product replacement model

This commit is contained in:
QiuSW
2026-08-28 17:33:15 +08:00
parent 30d823880c
commit 696e0bca52
14 changed files with 1026 additions and 4 deletions
+2
View File
@@ -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)
+2
View File
@@ -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},
@@ -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)
}
}
}
+2
View File
@@ -41,6 +41,8 @@ func MigratedModels() []any {
&models.CollectionColorPrice{},
&models.CollectionSKU{},
&models.CollectionSKUValue{},
&models.PDDProductReplacement{},
&models.PDDProductReplacementItem{},
}
}
@@ -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 {
+126
View File
@@ -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
}
+115
View File
@@ -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
}
@@ -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)
}
+15
View File
@@ -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)
}
+335
View File
@@ -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}
}
@@ -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
}
@@ -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
})
}