feat(#131): activate product replacements safely

This commit is contained in:
QiuSW
2026-08-28 17:56:17 +08:00
parent 696e0bca52
commit 92385c4a8d
17 changed files with 1458 additions and 73 deletions
+12 -2
View File
@@ -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: d585a906857a46752788c9f25d3ac0361739dada
synchronized_at: 2026-08-28T09:27:21Z
wiki_revision: b54f4cce004cb6eb2398b40f89ad88888b773a3a
synchronized_at: 2026-08-28T09:50:04Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
@@ -234,3 +234,13 @@ Android Portal/Agent
- 主表的可空 `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 权限。
## PDD 商品替换生效与规格匹配(#131)
- `replacement.Service.RegisterAndActivate` 在同一数据库事务内锁定受影响采购任务与虾皮商品:虾皮商品改指替代 PDD、清除旧规格映射,待执行/待探测采购任务取消并释放租约与运行守卫;执行中和终态任务的 PDD 外键及全部快照保持原样,旧 PDD 仅置为 `disabled`。
- `pdd_product_replacement_item` 是持久匹配工作项;`pdd_product_replacement_worker_lease` 为多实例全局租约。事务提交后异步唤醒,API 服务启动时恢复遗留 `matching`,最多有限重试,失败转 `manual_required`。
- 匹配复用 `aimatching.Service.Resolve`。确定性 `exact_match` 可直接确认且不要求置信度;`ai_match` 只有在开关启用、置信度达到 `ai_matching_setting.auto_confirm_min_confidence`(默认 0.9)、候选/角色/原因/输入版本全部通过校验时才确认。
- 人工维护虾皮规格映射后,同一事务同步分项状态与主表 `completed` / `completed_partial` 聚合;worker 只更新仍为 `matching` 且商品关联、规格 JSON 未变化的行,不覆盖人工结果。
- 纠错使用 `CorrectAndActivate`:原记录转 `superseded`,新建原始失效商品到新目标的记录,只迁移原 `replacement_item` 冻结的虾皮商品集合,不影响共享中间目标的其他商品。
- 历史采购执行身份继续以 `PDDGoodsIDSnapshot` / `PDDURLSnapshot` 为准;当前采购查询未发现按可变 `pdd_product_id` 聚合历史数据的实现,因此本工单无需改写统计 SQL。
+13 -2
View File
@@ -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: b31ebf8888da7c3d9b71520938b104ff9e66511f
synchronized_at: 2026-08-28T09:27:33Z
wiki_revision: 14d111bf1043f9c31f4f54dbfbca38826e986ab9
synchronized_at: 2026-08-28T09:50:14Z
<!-- gitea-wiki-mirror:end -->
# 业务规则与术语
@@ -307,3 +307,14 @@ synchronized_at: 2026-08-28T09:27:33Z
- 选错替代商品时,连续执行 B→C 不等于撤销 A→B,因为 B 可能还关联其他虾皮商品。正确纠错语义是把原 A→B 记录置为 `superseded`,再建立 A→C,并只处理原记录分项中冻结的影响集合。
- 创建请求按 `create_request_id` 幂等;重放时源商品、替代商品、来源类型、来源任务、采集证据和发起设备必须一致,否则返回幂等冲突,不能静默覆盖。
- `created_by_device_id` 只代表 Agent 设备。当前系统没有设备到采购员账号的绑定,多人多机场景若需要个人责任追踪,必须另建工单实现设备绑定操作员。
## PDD 商品替换生效与规格匹配(#131)
- 替代商品采集成功后,服务端自动、原子生效,无人工审批。受影响虾皮商品改指替代商品并清除旧规格映射;旧 PDD 商品置为停用但不删除。
- 原有 `pending` / `spec_probe_pending` 采购任务取消并标记 `PDD_PRODUCT_REPLACED`,释放租约与运行守卫,之后只能由既有重试流程按最新档案创建新任务;不得把新商品外键与旧 URL、goods_id、价格或规格快照拼在同一任务中。
- `running` / `order_submit_started` 与所有终态采购任务继续保留原商品外键和全部快照,按原任务事实收敛或留作历史;`TargetColorSnapshot` / `TargetSizeSnapshot` 始终是虾皮需求,不因商品替换而清除。
- 规格匹配按虾皮商品分项持久执行。AI 开关关闭、无结果或未通过自动确认门槛时转人工;确定性唯一 `exact_match` 可不带置信度自动确认。
- 用户于 2026-08-28 确认放宽 #46 的 AI 口径:`ai_match` 仅在置信度存在且不低于服务端阈值(默认 0.9)、结果严格属于当前候选、颜色/尺码角色唯一完整、原因非空且写入前输入版本未变化时自动确认。否则不得猜测或选择相近候选,转 `manual_required`。
- 自动确认必须保存来源、实际置信度和脱敏限长原因;人工确认完成后同步对应分项及主表总体状态。worker 不得覆盖已由人工改变的映射。
- 纠错不能简单执行 B→C;必须把原 A→B 记录置为 `superseded`,建立 A→C,并只处理原记录冻结的虾皮商品影响集合。
+15 -8
View File
@@ -17,11 +17,12 @@ import (
)
const (
settingID = uint64(1)
defaultTimeoutSecs = 15
maxTimeoutSecs = 600
minTimeoutSecs = 3
maxSettingText = 2048
settingID = uint64(1)
defaultTimeoutSecs = 15
maxTimeoutSecs = 600
minTimeoutSecs = 3
maxSettingText = 2048
defaultAutoConfirmMinConfidence = 0.9
)
// Service owns the internal AI Provider configuration. The API key exception
@@ -45,7 +46,7 @@ func secureHTTPClient() *http.Client {
func (s *Service) Settings(ctx context.Context) (SettingsView, error) {
setting, err := s.setting(ctx)
if errors.Is(err, gorm.ErrRecordNotFound) {
return SettingsView{Enabled: false, Provider: ProviderOpenAICompatible, TimeoutSeconds: defaultTimeoutSecs}, nil
return SettingsView{Enabled: false, Provider: ProviderOpenAICompatible, TimeoutSeconds: defaultTimeoutSecs, AutoConfirmMinConfidence: defaultAutoConfirmMinConfidence}, nil
}
if err != nil {
return SettingsView{}, err
@@ -69,7 +70,7 @@ func (s *Service) SaveSettings(ctx context.Context, request SaveSettingsRequest,
return err
}
if errors.Is(err, gorm.ErrRecordNotFound) {
current = models.AIMatchingSetting{ID: settingID, Provider: ProviderOpenAICompatible, BaseURL: "", Model: "", APIKey: "", TimeoutSeconds: defaultTimeoutSecs}
current = models.AIMatchingSetting{ID: settingID, Provider: ProviderOpenAICompatible, BaseURL: "", Model: "", APIKey: "", TimeoutSeconds: defaultTimeoutSecs, AutoConfirmMinConfidence: defaultAutoConfirmMinConfidence}
}
if key := strings.TrimSpace(request.APIKey); key != "" {
current.APIKey = key
@@ -82,6 +83,9 @@ func (s *Service) SaveSettings(ctx context.Context, request SaveSettingsRequest,
current.BaseURL = request.BaseURL
current.Model = request.Model
current.TimeoutSeconds = request.TimeoutSeconds
if request.AutoConfirmMinConfidence != nil {
current.AutoConfirmMinConfidence = *request.AutoConfirmMinConfidence
}
current.UpdatedBy = &operatorID
if errors.Is(err, gorm.ErrRecordNotFound) {
if createErr := tx.Create(&current).Error; createErr != nil {
@@ -208,6 +212,9 @@ func validateSettings(request SaveSettingsRequest) (SaveSettingsRequest, error)
if len(request.BaseURL) > maxSettingText || len(request.Model) > 255 || len(request.APIKey) > maxSettingText {
return SaveSettingsRequest{}, fail(CodeInvalidSetting, "AI 设置内容过长")
}
if request.AutoConfirmMinConfidence != nil && (*request.AutoConfirmMinConfidence < 0 || *request.AutoConfirmMinConfidence > 1) {
return SaveSettingsRequest{}, fail(CodeInvalidSetting, "自动确认阈值必须在 0 与 1 之间")
}
if request.Enabled {
if request.BaseURL == "" || request.Model == "" {
return SaveSettingsRequest{}, fail(CodeInvalidSetting, "启用 AI 匹配前请填写服务地址和模型")
@@ -236,7 +243,7 @@ func settingView(setting models.AIMatchingSetting) SettingsView {
if timeout == 0 {
timeout = defaultTimeoutSecs
}
return SettingsView{Enabled: setting.Enabled, Provider: ProviderOpenAICompatible, BaseURL: setting.BaseURL, Model: setting.Model, APIKey: setting.APIKey, TimeoutSeconds: timeout}
return SettingsView{Enabled: setting.Enabled, Provider: ProviderOpenAICompatible, BaseURL: setting.BaseURL, Model: setting.Model, APIKey: setting.APIKey, TimeoutSeconds: timeout, AutoConfirmMinConfidence: setting.AutoConfirmMinConfidence}
}
func (s *Service) httpClient() *http.Client {
@@ -113,3 +113,25 @@ func TestSaveSettingsAllowsTimeoutUpToSixHundredSeconds(t *testing.T) {
t.Fatalf("601-second timeout must be rejected: %v", err)
}
}
func TestAutoConfirmThresholdDefaultsPersistsAndValidates(t *testing.T) {
service := matcherTestService(t)
view, err := service.Settings(context.Background())
if err != nil || view.AutoConfirmMinConfidence != 0.9 {
t.Fatalf("default view=%+v err=%v", view, err)
}
threshold := 0.95
view, err = service.SaveSettings(context.Background(), SaveSettingsRequest{
Enabled: true, BaseURL: "https://example.com/v1", Model: "test", APIKey: "test-key", TimeoutSeconds: 15,
AutoConfirmMinConfidence: &threshold,
}, 1)
if err != nil || view.AutoConfirmMinConfidence != threshold {
t.Fatalf("saved view=%+v err=%v", view, err)
}
invalid := 1.01
_, err = service.SaveSettings(context.Background(), SaveSettingsRequest{TimeoutSeconds: 15, AutoConfirmMinConfidence: &invalid}, 1)
var settingErr *Error
if !errors.As(err, &settingErr) || settingErr.Code != CodeInvalidSetting {
t.Fatalf("invalid threshold err=%v", err)
}
}
+13 -11
View File
@@ -55,20 +55,22 @@ type MatchResult struct {
}
type SettingsView struct {
Enabled bool `json:"enabled"`
Provider string `json:"provider,omitempty"`
BaseURL string `json:"baseUrl,omitempty"`
Model string `json:"model,omitempty"`
TimeoutSeconds int `json:"timeoutSeconds,omitempty"`
APIKey string `json:"apiKey,omitempty"`
Enabled bool `json:"enabled"`
Provider string `json:"provider,omitempty"`
BaseURL string `json:"baseUrl,omitempty"`
Model string `json:"model,omitempty"`
TimeoutSeconds int `json:"timeoutSeconds,omitempty"`
AutoConfirmMinConfidence float64 `json:"autoConfirmMinConfidence"`
APIKey string `json:"apiKey,omitempty"`
}
type SaveSettingsRequest struct {
Enabled bool `json:"enabled"`
BaseURL string `json:"baseUrl"`
Model string `json:"model"`
APIKey string `json:"apiKey"`
TimeoutSeconds int `json:"timeoutSeconds"`
Enabled bool `json:"enabled"`
BaseURL string `json:"baseUrl"`
Model string `json:"model"`
APIKey string `json:"apiKey"`
TimeoutSeconds int `json:"timeoutSeconds"`
AutoConfirmMinConfidence *float64 `json:"autoConfirmMinConfidence,omitempty"`
}
type Error struct {
+1
View File
@@ -43,6 +43,7 @@ func MigratedModels() []any {
&models.CollectionSKUValue{},
&models.PDDProductReplacement{},
&models.PDDProductReplacementItem{},
&models.PDDProductReplacementWorkerLease{},
}
}
+19 -3
View File
@@ -51,9 +51,10 @@ type PDDProductReplacement struct {
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"`
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"`
ActivatedAt *time.Time `json:"activatedAt,omitempty" gorm:"index"`
}
func (PDDProductReplacement) TableName() string { return "pdd_product_replacement" }
@@ -124,3 +125,18 @@ func (item *PDDProductReplacementItem) BeforeCreate(_ *gorm.DB) error {
}
return nil
}
// PDDProductReplacementWorkerLease serializes persistent matching workers
// across service instances. Work remains in replacement_item; the lease only
// elects one owner and is safe to steal after expiry.
type PDDProductReplacementWorkerLease struct {
ID uint8 `json:"-" gorm:"primaryKey;autoIncrement:false"`
OwnerID string `json:"-" gorm:"size:64;not null;default:''"`
ExpiresAt time.Time `json:"-" gorm:"not null;index"`
Version uint64 `json:"-" gorm:"not null;default:0"`
UpdatedAt time.Time `json:"-"`
}
func (PDDProductReplacementWorkerLease) TableName() string {
return "pdd_product_replacement_worker_lease"
}
+11 -10
View File
@@ -73,16 +73,17 @@ func (PDDProduct) TableName() string { return "pdd_product" }
// administrator settings endpoint. It must never be added to task records,
// Agent payloads, purchaser responses, logs, code, tickets, or Wiki pages.
type AIMatchingSetting struct {
ID uint64 `json:"id" gorm:"primaryKey;autoIncrement:false"`
Enabled bool `json:"enabled" gorm:"not null;default:false"`
Provider string `json:"provider" gorm:"size:32;not null;default:openai_compatible"`
BaseURL string `json:"baseUrl" gorm:"type:text;not null"`
Model string `json:"model" gorm:"size:255;not null;default:''"`
APIKey string `json:"apiKey" gorm:"column:api_key;type:text;not null"`
TimeoutSeconds int `json:"timeoutSeconds" gorm:"not null;default:15"`
UpdatedBy *uint64 `json:"updatedBy"`
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
ID uint64 `json:"id" gorm:"primaryKey;autoIncrement:false"`
Enabled bool `json:"enabled" gorm:"not null;default:false"`
Provider string `json:"provider" gorm:"size:32;not null;default:openai_compatible"`
BaseURL string `json:"baseUrl" gorm:"type:text;not null"`
Model string `json:"model" gorm:"size:255;not null;default:''"`
APIKey string `json:"apiKey" gorm:"column:api_key;type:text;not null"`
TimeoutSeconds int `json:"timeoutSeconds" gorm:"not null;default:15"`
AutoConfirmMinConfidence float64 `json:"autoConfirmMinConfidence" gorm:"not null;default:0.9;check:ck_ai_matching_auto_confirm_confidence,auto_confirm_min_confidence >= 0 AND auto_confirm_min_confidence <= 1"`
UpdatedBy *uint64 `json:"updatedBy"`
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
}
func (AIMatchingSetting) TableName() string { return "ai_matching_setting" }
+378
View File
@@ -0,0 +1,378 @@
package replacement
import (
"context"
"errors"
"time"
"go-admin/app/goauto/models"
"go-admin/app/goauto/shopeeproduct"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
const (
CodeNoAffectedShopeeProducts = "REPLACEMENT_NO_AFFECTED_SHOPEE_PRODUCTS"
CodeActivationConflict = "REPLACEMENT_ACTIVATION_CONFLICT"
PurchaseReplacedErrorCode = "PDD_PRODUCT_REPLACED"
)
type ActivationResult struct {
Record Record `json:"record"`
AffectedShopeeProducts int `json:"affectedShopeeProducts"`
CancelledPendingTasks int `json:"cancelledPendingTasks"`
SkippedRunningTasks int `json:"skippedRunningTasks"`
PreservedTerminalTasks int `json:"preservedTerminalTasks"`
Replayed bool `json:"replayed"`
}
// RegisterAndActivate is the production entry used by #130. Registration,
// Shopee relinking, mapping clearing, pending-task cancellation, source
// disabling and persistent work-item creation commit together or not at all.
func (service *Service) RegisterAndActivate(ctx context.Context, request RegisterRequest) (ActivationResult, error) {
if err := validateRegisterRequest(request); err != nil {
return ActivationResult{}, err
}
if service.DB == nil {
return ActivationResult{}, internal(errors.New("database is nil"))
}
var result ActivationResult
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
record, replayed, err := service.registerForActivation(tx, request)
if err != nil {
return err
}
result.Replayed = replayed
if record.ActivatedAt == nil {
if err := service.activate(tx, &record, request, &result); err != nil {
return err
}
} else if err := loadActivationCounts(tx, record, &result); err != nil {
return err
}
items := make([]models.PDDProductReplacementItem, 0)
if err := tx.Where("replacement_id = ?", record.ID).Order("id").Find(&items).Error; err != nil {
return internal(err)
}
result.Record = Record{Replacement: record, Items: items, Replayed: replayed}
return nil
})
if err == nil && service.StartMatching != nil {
service.StartMatching(service.DB)
}
return result, err
}
// CorrectAndActivate supersedes one activated replacement and creates the
// corrected original-source -> new-target fact. Only Shopee products frozen in
// the superseded record's items are moved; unrelated products already sharing
// the intermediate target are deliberately untouched.
func (service *Service) CorrectAndActivate(ctx context.Context, supersededID uint64, request RegisterRequest) (ActivationResult, error) {
if err := validateRegisterRequest(request); err != nil {
return ActivationResult{}, err
}
if supersededID == 0 || service.DB == nil {
return ActivationResult{}, fail(CodeInvalidRequest, "supersededReplacementId 无效")
}
var result ActivationResult
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var previous models.PDDProductReplacement
if err := tx.First(&previous, supersededID).Error; errors.Is(err, gorm.ErrRecordNotFound) {
return fail(CodeNotFound, "待纠错替换记录不存在")
} else if err != nil {
return internal(err)
}
if request.SourceProductID != previous.TargetProductID {
return fail(CodeInvalidRequest, "纠错来源必须是原替代商品")
}
if previous.SourceProductID == request.TargetProductID {
return fail(CodeCycleDetected, "纠错不能把商品替换回原失效商品")
}
var replay models.PDDProductReplacement
if err := tx.Where("create_request_id = ?", request.RequestID).First(&replay).Error; err == nil {
if replay.SourceProductID != previous.SourceProductID || replay.TargetProductID != request.TargetProductID || replay.OriginType != request.OriginType || replay.OriginTaskID != request.OriginTaskID || replay.TargetCollectionTaskID != request.TargetCollectionTaskID || replay.CreatedByDeviceID != request.CreatedByDeviceID {
return fail(CodeIdempotencyConflict, "requestId 已用于其他替换请求")
}
result.Replayed = true
if err := loadActivationCounts(tx, replay, &result); err != nil {
return err
}
var items []models.PDDProductReplacementItem
if err := tx.Where("replacement_id = ?", replay.ID).Order("id").Find(&items).Error; err != nil {
return internal(err)
}
result.Record = Record{Replacement: replay, Items: items, Replayed: true}
return nil
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
return internal(err)
}
if previous.Status != models.ReplacementStatusActive || previous.ActivatedAt == nil {
return fail(CodeActivationConflict, "原替换记录已被纠错或尚未生效")
}
if err := validateOriginProduct(tx, request, previous.TargetProductID); err != nil {
return err
}
if err := validateCollectionProof(tx, request); err != nil {
return err
}
var target models.PDDProduct
if err := tx.First(&target, request.TargetProductID).Error; errors.Is(err, gorm.ErrRecordNotFound) {
return fail(CodeProductNotFound, "替代商品不存在")
} else if err != nil {
return internal(err)
} else if target.Status != "active" {
return fail(CodeTargetNotActive, "替代商品不是可用状态")
}
var previousItems []models.PDDProductReplacementItem
if err := tx.Where("replacement_id = ?", previous.ID).Order("id").Find(&previousItems).Error; err != nil {
return internal(err)
}
if len(previousItems) == 0 {
return fail(CodeNoAffectedShopeeProducts, "原替换记录没有影响范围")
}
shopeeIDs := make([]uint64, 0, len(previousItems))
for _, item := range previousItems {
shopeeIDs = append(shopeeIDs, item.ShopeeProductID)
}
var tasks []models.PurchaseTask
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("shopee_product_id IN ? AND pdd_product_id = ?", shopeeIDs, previous.TargetProductID).Order("id").Find(&tasks).Error; err != nil {
return internal(err)
}
var affected []models.ShopeeProduct
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id IN ? AND pdd_product_id = ?", shopeeIDs, previous.TargetProductID).Order("id").Find(&affected).Error; err != nil {
return internal(err)
}
if len(affected) != len(previousItems) {
return fail(CodeActivationConflict, "原替换影响范围已变化")
}
now := service.Now()
if update := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PDDProductReplacement{}).Where("id = ? AND status = ?", previous.ID, models.ReplacementStatusActive).Updates(map[string]any{"status": models.ReplacementStatusSuperseded, "active_slot": nil}); update.Error != nil {
return internal(update.Error)
} else if update.RowsAffected != 1 {
return fail(CodeActivationConflict, "原替换记录已变化")
}
corrected := models.PDDProductReplacement{
SourceProductID: previous.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, ActivatedAt: &now,
}
if err := tx.Create(&corrected).Error; err != nil {
if isUniqueConstraint(err) {
return fail(CodeSourceAlreadyReplaced, "该失效商品已有生效中的替换记录")
}
return internal(err)
}
if err := applyPurchaseReplacement(tx, tasks, now, &result); err != nil {
return err
}
items := make([]models.PDDProductReplacementItem, 0, len(affected))
for index := range affected {
product := &affected[index]
cleared, err := clearMappings(product.SpecsJSON)
if err != nil {
return internal(err)
}
if update := tx.Model(&models.ShopeeProduct{}).Where("id = ? AND pdd_product_id = ?", product.ID, previous.TargetProductID).Updates(map[string]any{"pdd_product_id": request.TargetProductID, "specs_json": cleared}); update.Error != nil {
return internal(update.Error)
} else if update.RowsAffected != 1 {
return fail(CodeActivationConflict, "虾皮商品关联已变化")
}
items = append(items, models.PDDProductReplacementItem{ReplacementID: corrected.ID, ShopeeProductID: product.ID, MappingStatus: models.ReplacementItemMappingMatching, UpdatedAt: now})
}
if err := tx.Create(&items).Error; err != nil {
return internal(err)
}
result.AffectedShopeeProducts = len(items)
result.Record = Record{Replacement: corrected, Items: items}
return nil
})
if err == nil && service.StartMatching != nil {
service.StartMatching(service.DB)
}
return result, err
}
func (service *Service) registerForActivation(tx *gorm.DB, request RegisterRequest) (models.PDDProductReplacement, bool, error) {
var existing models.PDDProductReplacement
if err := tx.Where("create_request_id = ?", request.RequestID).First(&existing).Error; err == nil {
if !sameRequest(existing, request) {
return models.PDDProductReplacement{}, false, fail(CodeIdempotencyConflict, "requestId 已用于其他替换请求")
}
return existing, true, nil
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
return models.PDDProductReplacement{}, false, internal(err)
}
if err := validateProducts(tx, request); err != nil {
return models.PDDProductReplacement{}, false, err
}
if err := validateOrigin(tx, request); err != nil {
return models.PDDProductReplacement{}, false, err
}
if err := validateCollectionProof(tx, request); err != nil {
return models.PDDProductReplacement{}, false, err
}
if err := ensureNoActiveReplacement(tx, request.SourceProductID); err != nil {
return models.PDDProductReplacement{}, false, err
}
if err := ensureNoCycle(tx, request.SourceProductID, request.TargetProductID); err != nil {
return models.PDDProductReplacement{}, false, 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 models.PDDProductReplacement{}, false, fail(CodeSourceAlreadyReplaced, "该失效商品已有生效中的替换记录")
}
return models.PDDProductReplacement{}, false, internal(err)
}
return record, false, nil
}
func (service *Service) activate(tx *gorm.DB, record *models.PDDProductReplacement, request RegisterRequest, result *ActivationResult) error {
// Read the affected Shopee IDs first, then lock purchase rows in ID order.
// Claim/Start also lock one purchase row first; this preserves that order and
// avoids a product->task deadlock.
var shopeeIDs []uint64
if err := tx.Model(&models.ShopeeProduct{}).Where("pdd_product_id = ?", request.SourceProductID).Order("id").Pluck("id", &shopeeIDs).Error; err != nil {
return internal(err)
}
if len(shopeeIDs) == 0 {
return fail(CodeNoAffectedShopeeProducts, "失效商品没有关联的虾皮商品")
}
var tasks []models.PurchaseTask
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("shopee_product_id IN ? AND pdd_product_id = ?", shopeeIDs, request.SourceProductID).Order("id").Find(&tasks).Error; err != nil {
return internal(err)
}
var affected []models.ShopeeProduct
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id IN ? AND pdd_product_id = ?", shopeeIDs, request.SourceProductID).Order("id").Find(&affected).Error; err != nil {
return internal(err)
}
if len(affected) != len(shopeeIDs) {
return fail(CodeActivationConflict, "虾皮商品关联已变化,请刷新后重试")
}
if err := validateProducts(tx, request); err != nil {
return err
}
if err := validateCollectionProof(tx, request); err != nil {
return err
}
now := service.Now()
if err := applyPurchaseReplacement(tx, tasks, now, result); err != nil {
return err
}
items := make([]models.PDDProductReplacementItem, 0, len(affected))
for index := range affected {
product := &affected[index]
cleared, err := clearMappings(product.SpecsJSON)
if err != nil {
return internal(err)
}
update := tx.Model(&models.ShopeeProduct{}).Where("id = ? AND pdd_product_id = ?", product.ID, request.SourceProductID).
Updates(map[string]any{"pdd_product_id": request.TargetProductID, "specs_json": cleared})
if update.Error != nil {
return internal(update.Error)
}
if update.RowsAffected != 1 {
return fail(CodeActivationConflict, "虾皮商品关联已变化,请重试")
}
items = append(items, models.PDDProductReplacementItem{ReplacementID: record.ID, ShopeeProductID: product.ID, MappingStatus: models.ReplacementItemMappingMatching, UpdatedAt: now})
}
if err := tx.Create(&items).Error; err != nil {
return internal(err)
}
if update := tx.Model(&models.PDDProduct{}).Where("id = ?", request.SourceProductID).Update("status", "disabled"); update.Error != nil {
return internal(update.Error)
} else if update.RowsAffected != 1 {
return fail(CodeActivationConflict, "失效商品状态已变化,请重试")
}
record.ActivatedAt = &now
record.MappingUpdatedAt = now
if err := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PDDProductReplacement{}).Where("id = ? AND activated_at IS NULL", record.ID).
Updates(map[string]any{"activated_at": now, "mapping_updated_at": now, "mapping_status": models.ReplacementMappingMatching}).Error; err != nil {
return internal(err)
}
result.AffectedShopeeProducts = len(affected)
return nil
}
func applyPurchaseReplacement(tx *gorm.DB, tasks []models.PurchaseTask, now time.Time, result *ActivationResult) error {
for index := range tasks {
task := &tasks[index]
switch task.Status {
case models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending:
reason := "PDD 商品已替换,请在规格匹配完成后继续采购"
update := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).
Where("id = ? AND status IN ?", task.ID, []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending}).
Updates(map[string]any{
"status": models.PurchaseTaskStatusCancelled, "status_version": gorm.Expr("status_version + 1"), "status_changed_at": now,
"active_slot": nil, "device_run_slot": nil, "account_run_slot": nil, "lease_expires_at": nil,
"cancelled_at": now, "cancel_reason": reason, "error_code": PurchaseReplacedErrorCode, "error_message": reason,
})
if update.Error != nil {
return internal(update.Error)
}
if update.RowsAffected != 1 {
return fail(CodeActivationConflict, "采购任务状态已变化,请重试")
}
result.CancelledPendingTasks++
case models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted:
result.SkippedRunningTasks++
default:
result.PreservedTerminalTasks++
}
}
return nil
}
func clearMappings(raw string) (string, error) {
specs, err := shopeeproduct.Unmarshal(raw)
if err != nil {
return "", err
}
for dimensionIndex := range specs {
for valueIndex := range specs[dimensionIndex].Values {
specs[dimensionIndex].Values[valueIndex].Mapping = nil
}
}
return shopeeproduct.Marshal(specs)
}
func loadActivationCounts(tx *gorm.DB, record models.PDDProductReplacement, result *ActivationResult) error {
var items []models.PDDProductReplacementItem
if err := tx.Where("replacement_id = ?", record.ID).Find(&items).Error; err != nil {
return internal(err)
}
result.AffectedShopeeProducts = len(items)
if len(items) == 0 {
return nil
}
ids := make([]uint64, 0, len(items))
for _, item := range items {
ids = append(ids, item.ShopeeProductID)
}
var tasks []models.PurchaseTask
if err := tx.Where("shopee_product_id IN ? AND pdd_product_id = ?", ids, record.SourceProductID).Find(&tasks).Error; err != nil {
return internal(err)
}
for _, task := range tasks {
if task.Status == models.PurchaseTaskStatusCancelled && task.ErrorCode != nil && *task.ErrorCode == PurchaseReplacedErrorCode {
result.CancelledPendingTasks++
} else if task.Status == models.PurchaseTaskStatusRunning || task.Status == models.PurchaseTaskStatusOrderSubmitStarted {
result.SkippedRunningTasks++
} else {
result.PreservedTerminalTasks++
}
}
return nil
}
@@ -0,0 +1,185 @@
package replacement
import (
"context"
"fmt"
"testing"
"time"
"go-admin/app/goauto/models"
"github.com/google/uuid"
"gorm.io/gorm"
)
func TestRegisterAndActivateAppliesPurchaseStateMatrixAtomically(t *testing.T) {
fixture := seedReplacementFixture(t)
fixture.service.StartMatching = func(*gorm.DB) {}
if err := fixture.db.Model(&models.PDDProduct{}).Where("id = ?", fixture.source.ID).Update("status", "active").Error; err != nil {
t.Fatal(err)
}
statuses := []string{
models.PurchaseTaskStatusPending,
models.PurchaseTaskStatusSpecProbePending,
models.PurchaseTaskStatusRunning,
models.PurchaseTaskStatusOrderSubmitStarted,
models.PurchaseTaskStatusRehearsalCompleted,
models.PurchaseTaskStatusOrderCreated,
models.PurchaseTaskStatusOrderResultUnknown,
models.PurchaseTaskStatusFailed,
models.PurchaseTaskStatusCancelled,
}
tasks := make([]models.PurchaseTask, 0, len(statuses))
shopees := make([]models.ShopeeProduct, 0, len(statuses))
for index, status := range statuses {
shopee := models.ShopeeProduct{
ShopeeItemID: fmt.Sprintf("ACT-%d", index), PDDProductID: &fixture.source.ID, Currency: "CNY",
SpecsJSON: `[{"name":"颜色","role":"color","values":[{"name":"红色","source":"import","mapping":{"pddValue":"旧红色","source":"manual","status":"confirmed"}}]}]`,
}
if err := fixture.db.Create(&shopee).Error; err != nil {
t.Fatal(err)
}
shopees = append(shopees, shopee)
task := activationPurchaseTask(fixture, shopee, status, index)
if err := fixture.db.Session(&gorm.Session{SkipHooks: true}).Create(&task).Error; err != nil {
t.Fatalf("create %s: %v", status, err)
}
tasks = append(tasks, task)
}
result, err := fixture.service.RegisterAndActivate(context.Background(), fixture.request())
if err != nil {
t.Fatal(err)
}
if result.AffectedShopeeProducts != len(statuses) || result.CancelledPendingTasks != 2 || result.SkippedRunningTasks != 2 || result.PreservedTerminalTasks != 5 {
t.Fatalf("unexpected counts: %+v", result)
}
for index, original := range tasks {
var got models.PurchaseTask
if err := fixture.db.First(&got, original.ID).Error; err != nil {
t.Fatal(err)
}
if got.PDDProductID != fixture.source.ID || got.PDDGoodsIDSnapshot != original.PDDGoodsIDSnapshot || got.PDDURLSnapshot != original.PDDURLSnapshot {
t.Fatalf("task %s mixed replacement identity with old snapshots: %+v", statuses[index], got)
}
if index < 2 {
if got.Status != models.PurchaseTaskStatusCancelled || got.ErrorCode == nil || *got.ErrorCode != PurchaseReplacedErrorCode || got.ActiveSlot != nil || got.DeviceRunSlot != nil || got.AccountRunSlot != nil || got.LeaseExpiresAt != nil {
t.Fatalf("pending task not safely cancelled: %+v", got)
}
} else if got.Status != original.Status || got.StatusVersion != original.StatusVersion {
t.Fatalf("protected task changed: before=%+v after=%+v", original, got)
}
}
for _, original := range shopees {
var got models.ShopeeProduct
if err := fixture.db.First(&got, original.ID).Error; err != nil {
t.Fatal(err)
}
if got.PDDProductID == nil || *got.PDDProductID != fixture.target.ID || got.SpecsJSON == original.SpecsJSON {
t.Fatalf("Shopee product was not relinked and cleared: %+v", got)
}
}
var source models.PDDProduct
if err := fixture.db.First(&source, fixture.source.ID).Error; err != nil || source.Status != "disabled" {
t.Fatalf("source status=%s err=%v", source.Status, err)
}
var itemCount int64
if err := fixture.db.Model(&models.PDDProductReplacementItem{}).Where("replacement_id = ?", result.Record.Replacement.ID).Count(&itemCount).Error; err != nil || itemCount != int64(len(statuses)) {
t.Fatalf("item count=%d err=%v", itemCount, err)
}
replay, err := fixture.service.RegisterAndActivate(context.Background(), fixture.requestWithID(result.Record.Replacement.CreateRequestID))
if err != nil || !replay.Replayed || replay.Record.Replacement.ID != result.Record.Replacement.ID || replay.AffectedShopeeProducts != len(statuses) {
t.Fatalf("replay=%+v err=%v", replay, err)
}
}
func TestRegisterAndActivateRollsBackWhenNoShopeeProductsAreAffected(t *testing.T) {
fixture := seedReplacementFixture(t)
fixture.service.StartMatching = func(*gorm.DB) { t.Fatal("worker started after rollback") }
_, err := fixture.service.RegisterAndActivate(context.Background(), fixture.request())
if replacementCode(err) != CodeNoAffectedShopeeProducts {
t.Fatalf("code=%s err=%v", replacementCode(err), err)
}
var count int64
if err := fixture.db.Model(&models.PDDProductReplacement{}).Count(&count).Error; err != nil || count != 0 {
t.Fatalf("replacement count=%d err=%v", count, err)
}
}
func TestCorrectAndActivateUsesOnlySupersededItemScope(t *testing.T) {
fixture := seedReplacementFixture(t)
fixture.service.StartMatching = func(*gorm.DB) {}
affected := models.ShopeeProduct{ShopeeItemID: "CORRECT-AFFECTED", PDDProductID: &fixture.source.ID, Currency: "CNY", SpecsJSON: "[]"}
if err := fixture.db.Create(&affected).Error; err != nil {
t.Fatal(err)
}
first, err := fixture.service.RegisterAndActivate(context.Background(), fixture.request())
if err != nil {
t.Fatal(err)
}
unrelated := models.ShopeeProduct{ShopeeItemID: "CORRECT-UNRELATED", PDDProductID: &fixture.target.ID, Currency: "CNY", SpecsJSON: "[]"}
correctedTarget := models.PDDProduct{GoodsID: "700000000004", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=700000000004", Status: "active", SpecsJSON: "[]"}
if err := fixture.db.Create(&unrelated).Error; err != nil {
t.Fatal(err)
}
if err := fixture.db.Create(&correctedTarget).Error; err != nil {
t.Fatal(err)
}
failed := models.CollectionTask{PDDProductID: &fixture.target.ID, RuleID: fixture.rule.ID, DeviceID: &fixture.device.ID, Source: models.CollectionTaskSourceAgentCurrentPage, Status: models.TaskStatusFailed, URLSnapshot: fixture.target.URL, GoodsIDSnapshot: fixture.target.GoodsID, RuleSnapshot: `{}`}
proof := models.CollectionTask{PDDProductID: &correctedTarget.ID, RuleID: fixture.rule.ID, DeviceID: &fixture.device.ID, Source: models.CollectionTaskSourceAgentCurrentPage, Status: models.TaskStatusCompletedPartial, URLSnapshot: correctedTarget.URL, GoodsIDSnapshot: correctedTarget.GoodsID, RuleSnapshot: `{}`}
if err := fixture.db.Create(&failed).Error; err != nil {
t.Fatal(err)
}
if err := fixture.db.Create(&proof).Error; err != nil {
t.Fatal(err)
}
request := RegisterRequest{RequestID: uuid.NewString(), SourceProductID: fixture.target.ID, TargetProductID: correctedTarget.ID, OriginType: models.ReplacementOriginCollection, OriginTaskID: failed.ID, TargetCollectionTaskID: proof.ID, CreatedByDeviceID: fixture.device.ID}
corrected, err := fixture.service.CorrectAndActivate(context.Background(), first.Record.Replacement.ID, request)
if err != nil {
t.Fatal(err)
}
if corrected.Record.Replacement.SourceProductID != fixture.source.ID || corrected.Record.Replacement.TargetProductID != correctedTarget.ID || corrected.AffectedShopeeProducts != 1 {
t.Fatalf("corrected=%+v", corrected)
}
var previous models.PDDProductReplacement
if err := fixture.db.First(&previous, first.Record.Replacement.ID).Error; err != nil || previous.Status != models.ReplacementStatusSuperseded || previous.ActiveSlot != nil {
t.Fatalf("previous=%+v err=%v", previous, err)
}
var gotAffected, gotUnrelated models.ShopeeProduct
_ = fixture.db.First(&gotAffected, affected.ID).Error
_ = fixture.db.First(&gotUnrelated, unrelated.ID).Error
if gotAffected.PDDProductID == nil || *gotAffected.PDDProductID != correctedTarget.ID {
t.Fatalf("affected target=%v", gotAffected.PDDProductID)
}
if gotUnrelated.PDDProductID == nil || *gotUnrelated.PDDProductID != fixture.target.ID {
t.Fatalf("unrelated product was moved: %v", gotUnrelated.PDDProductID)
}
}
func (fixture replacementFixture) requestWithID(requestID string) RegisterRequest {
request := fixture.request()
request.RequestID = requestID
return request
}
func activationPurchaseTask(fixture replacementFixture, shopee models.ShopeeProduct, status string, index int) models.PurchaseTask {
one := uint8(1)
now := time.Date(2026, 8, 28, 9, index, 0, 0, time.UTC)
lease := now.Add(time.Hour)
task := models.PurchaseTask{
ShopeeProductID: &shopee.ID, PDDProductID: fixture.source.ID, DeviceID: &fixture.device.ID,
ExecutionMode: models.PurchaseExecutionModeLive, Status: status,
ShopeeItemIDSnapshot: shopee.ShopeeItemID, PDDURLSnapshot: fixture.source.URL, PDDGoodsIDSnapshot: fixture.source.GoodsID,
PDDTitleSnapshot: "old-title", TargetColorSnapshot: "红色", MappedColorSnapshot: "旧红色", SpecSource: models.ReplacementItemSourceManualMapping,
SpecDecisionSnapshot: `{}`, Quantity: 1, Currency: "CNY", RuleType: "pddPurchase", RuleSchemaVersion: 1,
RequiredCapabilitiesJSON: `[]`, RuleSnapshot: `{}`, CreateRequestID: uuid.NewString(), StatusVersion: 7, StatusChangedAt: now,
PaymentReviewStatus: models.PurchasePaymentReviewPending, LogisticsStatus: models.PurchaseLogisticsStatusPending, WritebackStatus: models.PurchaseWritebackStatusNotSelected,
}
if status == models.PurchaseTaskStatusPending {
task.ActiveSlot, task.DeviceRunSlot, task.AccountRunSlot, task.LeaseExpiresAt = &one, &one, &one, &lease
} else if status == models.PurchaseTaskStatusSpecProbePending {
task.ActiveSlot, task.LeaseExpiresAt = &one, &lease
}
return task
}
+15 -7
View File
@@ -11,7 +11,6 @@ import (
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
const (
@@ -76,12 +75,17 @@ type ListResponse struct {
}
type Service struct {
DB *gorm.DB
Now func() time.Time
DB *gorm.DB
Now func() time.Time
StartMatching func(*gorm.DB)
}
func NewService(db *gorm.DB) *Service {
return &Service{DB: db, Now: func() time.Time { return time.Now().UTC() }}
return &Service{
DB: db,
Now: func() time.Time { return time.Now().UTC() },
StartMatching: StartMatching,
}
}
// Register writes the audit fact only. #131 owns the activation transaction
@@ -229,7 +233,7 @@ func validateRegisterRequest(request RegisterRequest) error {
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 {
if err := tx.Where("id IN ?", []uint64{request.SourceProductID, request.TargetProductID}).Find(&products).Error; err != nil {
return internal(err)
}
if len(products) != 2 {
@@ -244,6 +248,10 @@ func validateProducts(tx *gorm.DB, request RegisterRequest) error {
}
func validateOrigin(tx *gorm.DB, request RegisterRequest) error {
return validateOriginProduct(tx, request, request.SourceProductID)
}
func validateOriginProduct(tx *gorm.DB, request RegisterRequest, expectedProductID uint64) error {
switch request.OriginType {
case models.ReplacementOriginCollection:
var task models.CollectionTask
@@ -252,7 +260,7 @@ func validateOrigin(tx *gorm.DB, request RegisterRequest) error {
} 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 {
if task.Status != models.TaskStatusFailed || task.DeviceID == nil || *task.DeviceID != request.CreatedByDeviceID || task.PDDProductID == nil || *task.PDDProductID != expectedProductID {
return fail(CodeOriginTaskInvalid, "来源采集任务不满足替换条件")
}
case models.ReplacementOriginPurchase:
@@ -262,7 +270,7 @@ func validateOrigin(tx *gorm.DB, request RegisterRequest) error {
} else if err != nil {
return internal(err)
}
if task.Status != models.PurchaseTaskStatusFailed || task.DeviceID == nil || *task.DeviceID != request.CreatedByDeviceID || task.PDDProductID != request.SourceProductID {
if task.Status != models.PurchaseTaskStatusFailed || task.DeviceID == nil || *task.DeviceID != request.CreatedByDeviceID || task.PDDProductID != expectedProductID {
return fail(CodeOriginTaskInvalid, "来源采购任务不满足替换条件")
}
}
+424
View File
@@ -0,0 +1,424 @@
package replacement
import (
"context"
"encoding/json"
"errors"
"strings"
"time"
"go-admin/app/goauto/aimatching"
"go-admin/app/goauto/models"
"go-admin/app/goauto/shopeeproduct"
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
const (
defaultMatchAttempts = 3
defaultWorkerLease = 15 * time.Minute
)
type ReplacementMatcher interface {
Resolve(context.Context, aimatching.MatchRequest) (aimatching.MatchResult, error)
}
type Worker struct {
DB *gorm.DB
Matcher ReplacementMatcher
OwnerID string
Now func() time.Time
MaxAttempts int
LeaseDuration time.Duration
}
func NewWorker(db *gorm.DB) *Worker {
return &Worker{DB: db, Matcher: aimatching.NewService(db), OwnerID: uuid.NewString(), Now: func() time.Time { return time.Now().UTC() }, MaxAttempts: defaultMatchAttempts, LeaseDuration: defaultWorkerLease}
}
// StartMatching and RecoverMatching use the same persistent work queue. A
// process crash loses only the goroutine; matching rows remain claimable on the
// next trigger or service startup.
func StartMatching(db *gorm.DB) {
go func() {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
defer cancel()
_ = NewWorker(db).RunPending(ctx)
}()
}
func RecoverMatching(db *gorm.DB) { StartMatching(db) }
func (worker *Worker) RunPending(ctx context.Context) error {
for {
processed, err := worker.RunOnce(ctx)
if err != nil {
return err
}
if !processed {
return nil
}
}
}
func (worker *Worker) RunOnce(ctx context.Context) (bool, error) {
worker.defaults()
acquired, err := worker.acquire(ctx)
if err != nil || !acquired {
return false, err
}
defer worker.release()
var item models.PDDProductReplacementItem
err = worker.DB.WithContext(ctx).
Joins("JOIN pdd_product_replacement r ON r.id = pdd_product_replacement_item.replacement_id").
Where("pdd_product_replacement_item.mapping_status = ? AND pdd_product_replacement_item.attempt_count < ? AND r.status = ? AND r.activated_at IS NOT NULL", models.ReplacementItemMappingMatching, worker.MaxAttempts, models.ReplacementStatusActive).
Order("pdd_product_replacement_item.updated_at, pdd_product_replacement_item.id").
First(&item).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
return false, nil
}
if err != nil {
return false, internal(err)
}
now := worker.Now()
claimed := worker.DB.WithContext(ctx).Model(&models.PDDProductReplacementItem{}).
Where("id = ? AND mapping_status = ? AND attempt_count = ?", item.ID, models.ReplacementItemMappingMatching, item.AttemptCount).
Updates(map[string]any{"attempt_count": gorm.Expr("attempt_count + 1"), "updated_at": now})
if claimed.Error != nil {
return false, internal(claimed.Error)
}
if claimed.RowsAffected != 1 {
return true, nil
}
item.AttemptCount++
return true, worker.process(ctx, item)
}
func (worker *Worker) process(ctx context.Context, item models.PDDProductReplacementItem) error {
var replacement models.PDDProductReplacement
if err := worker.DB.WithContext(ctx).First(&replacement, item.ReplacementID).Error; err != nil {
return worker.finishFailure(ctx, item, "REPLACEMENT_NOT_FOUND", "替换记录不存在", false)
}
var shopee models.ShopeeProduct
if err := worker.DB.WithContext(ctx).First(&shopee, item.ShopeeProductID).Error; err != nil {
return worker.finishFailure(ctx, item, "SHOPEE_PRODUCT_NOT_FOUND", "虾皮商品不存在", false)
}
if shopee.PDDProductID == nil || *shopee.PDDProductID != replacement.TargetProductID {
return worker.finishFailure(ctx, item, "SHOPEE_LINK_CHANGED", "虾皮商品关联已变化", false)
}
var pdd models.PDDProduct
if err := worker.DB.WithContext(ctx).First(&pdd, replacement.TargetProductID).Error; err != nil || pdd.Status != "active" {
return worker.finishFailure(ctx, item, "PDD_PRODUCT_NOT_READY", "替代商品不可用", false)
}
settings, err := aimatching.NewService(worker.DB).Settings(ctx)
if err != nil {
return worker.finishFailure(ctx, item, "AI_SETTING_READ_FAILED", "读取 AI 匹配设置失败", true)
}
if !settings.Enabled {
return worker.finishFailure(ctx, item, aimatching.CodeNotConfigured, "AI 匹配未启用", false)
}
specs, err := shopeeproduct.Unmarshal(shopee.SpecsJSON)
if err != nil {
return worker.finishFailure(ctx, item, "SHOPEE_SPECS_INVALID", "虾皮规格数据无效", false)
}
candidates, err := replacementCandidates(pdd.SpecsJSON)
if err != nil {
return worker.finishFailure(ctx, item, "PDD_SPECS_INVALID", "PDD 规格数据无效", false)
}
planned := cloneSpecs(specs)
if !uniqueMatchRoles(planned) || candidates.colorRoles > 1 || candidates.sizeRoles > 1 {
return worker.finishFailure(ctx, item, "MATCH_ROLE_AMBIGUOUS", "颜色或尺码角色不唯一", false)
}
source := models.ReplacementItemSourceExactMatch
reasons := make([]string, 0)
var minimumConfidence *float64
relevant := 0
for dimensionIndex := range planned {
dimension := &planned[dimensionIndex]
if dimension.Role != shopeeproduct.RoleColor && dimension.Role != shopeeproduct.RoleSize {
continue
}
for valueIndex := range dimension.Values {
value := &dimension.Values[valueIndex]
relevant++
request := aimatching.MatchRequest{}
if dimension.Role == shopeeproduct.RoleColor {
request.TargetColor, request.Colors = value.Name, candidates.colors
} else {
request.TargetSize, request.Sizes = value.Name, candidates.sizes
}
matched, matchErr := worker.Matcher.Resolve(ctx, request)
if matchErr != nil {
return worker.finishFailure(ctx, item, matchErrorCode(matchErr), compactReason(matchErr.Error(), 500), retryableMatchError(matchErr))
}
pddValue := matched.MappedColor
if dimension.Role == shopeeproduct.RoleSize {
pddValue = matched.MappedSize
}
if !containsCandidate(pddValue, candidateValues(dimension.Role, candidates)) {
return worker.finishFailure(ctx, item, "MATCH_OUTSIDE_CANDIDATES", "匹配结果不在当前候选集中", false)
}
reason := strings.TrimSpace(matched.Decision.Reason)
mapping := &shopeeproduct.Mapping{PDDValue: pddValue, Source: matched.Source, Status: shopeeproduct.MappingStatusConfirmed, Reason: compactReason(reason, 500)}
switch matched.Source {
case aimatching.SourceExact:
if mapping.Reason == "" {
mapping.Reason = "规格名称标准化后唯一一致"
}
case aimatching.SourceAI:
confidence := matched.Decision.Confidence
if confidence == nil || *confidence < 0 || *confidence > 1 || *confidence < settings.AutoConfirmMinConfidence || mapping.Reason == "" {
return worker.finishFailure(ctx, item, "AI_AUTO_CONFIRM_REJECTED", "AI 结果未达到自动确认条件", false)
}
mapping.Confidence = confidence
source = models.ReplacementItemSourceAIMatch
if minimumConfidence == nil || *confidence < *minimumConfidence {
value := *confidence
minimumConfidence = &value
}
default:
return worker.finishFailure(ctx, item, "MATCH_SOURCE_INVALID", "匹配来源无效", false)
}
value.Mapping = mapping
reasons = append(reasons, mapping.Reason)
}
}
if relevant == 0 {
reasons = append(reasons, "商品没有需要映射的颜色或尺码")
}
if err := shopeeproduct.Validate(planned); err != nil {
return worker.finishFailure(ctx, item, "MATCH_RESULT_INVALID", compactReason(err.Error(), 500), false)
}
updatedJSON, err := shopeeproduct.Marshal(planned)
if err != nil {
return worker.finishFailure(ctx, item, "MATCH_RESULT_INVALID", "匹配结果无法保存", false)
}
return worker.finishMatched(ctx, item, replacement, shopee, updatedJSON, source, minimumConfidence, compactReason(strings.Join(reasons, ";"), 500))
}
func (worker *Worker) finishMatched(ctx context.Context, item models.PDDProductReplacementItem, replacement models.PDDProductReplacement, snapshot models.ShopeeProduct, specsJSON, source string, confidence *float64, reason string) error {
return worker.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var current models.PDDProductReplacementItem
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&current, item.ID).Error; err != nil {
return internal(err)
}
if current.MappingStatus != models.ReplacementItemMappingMatching {
return nil
}
productUpdate := tx.Model(&models.ShopeeProduct{}).
Where("id = ? AND pdd_product_id = ? AND specs_json = ?", snapshot.ID, replacement.TargetProductID, snapshot.SpecsJSON).
Update("specs_json", specsJSON)
if productUpdate.Error != nil {
return internal(productUpdate.Error)
}
if productUpdate.RowsAffected != 1 {
return worker.finishFailureTx(tx, current, "MATCH_INPUT_CHANGED", "规格或商品关联已变化", false)
}
now := worker.Now()
if err := tx.Model(&models.PDDProductReplacementItem{}).Where("id = ? AND mapping_status = ?", current.ID, models.ReplacementItemMappingMatching).
Updates(map[string]any{"mapping_status": models.ReplacementItemMappingMatched, "source": source, "confidence": confidence, "reason": reason, "last_error_code": nil, "last_error_at": nil, "updated_at": now}).Error; err != nil {
return internal(err)
}
return updateAggregate(tx, current.ReplacementID, now)
})
}
func (worker *Worker) finishFailure(ctx context.Context, item models.PDDProductReplacementItem, code, reason string, retryable bool) error {
return worker.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var current models.PDDProductReplacementItem
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&current, item.ID).Error; err != nil {
return internal(err)
}
if current.MappingStatus != models.ReplacementItemMappingMatching {
return nil
}
return worker.finishFailureTx(tx, current, code, reason, retryable)
})
}
func (worker *Worker) finishFailureTx(tx *gorm.DB, item models.PDDProductReplacementItem, code, reason string, retryable bool) error {
now := worker.Now()
updates := map[string]any{"last_error_code": compactReason(code, 64), "last_error_at": now, "reason": compactReason(reason, 500), "updated_at": now}
if !retryable || item.AttemptCount >= worker.MaxAttempts {
updates["mapping_status"] = models.ReplacementItemMappingManualRequired
}
if err := tx.Model(&models.PDDProductReplacementItem{}).Where("id = ? AND mapping_status = ?", item.ID, models.ReplacementItemMappingMatching).Updates(updates).Error; err != nil {
return internal(err)
}
return updateAggregate(tx, item.ReplacementID, now)
}
func updateAggregate(tx *gorm.DB, replacementID uint64, now time.Time) error {
var items []models.PDDProductReplacementItem
if err := tx.Where("replacement_id = ?", replacementID).Find(&items).Error; err != nil {
return internal(err)
}
status := models.ReplacementMappingMatching
if len(items) > 0 {
matched, terminal := 0, 0
for _, item := range items {
if item.MappingStatus == models.ReplacementItemMappingMatched {
matched++
}
if item.MappingStatus != models.ReplacementItemMappingMatching {
terminal++
}
}
if terminal == len(items) {
status = models.ReplacementMappingCompletedPartial
if matched == len(items) {
status = models.ReplacementMappingCompleted
}
}
}
return tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PDDProductReplacement{}).Where("id = ?", replacementID).Updates(map[string]any{"mapping_status": status, "mapping_updated_at": now}).Error
}
func (worker *Worker) acquire(ctx context.Context) (bool, error) {
now := worker.Now()
seed := models.PDDProductReplacementWorkerLease{ID: 1, OwnerID: "", ExpiresAt: time.Unix(0, 0).UTC()}
if err := worker.DB.WithContext(ctx).Clauses(clause.OnConflict{DoNothing: true}).Create(&seed).Error; err != nil {
return false, internal(err)
}
result := worker.DB.WithContext(ctx).Model(&models.PDDProductReplacementWorkerLease{}).
Where("id = 1 AND (expires_at < ? OR owner_id = ?)", now, worker.OwnerID).
Updates(map[string]any{"owner_id": worker.OwnerID, "expires_at": now.Add(worker.LeaseDuration), "version": gorm.Expr("version + 1")})
return result.RowsAffected == 1, result.Error
}
func (worker *Worker) release() {
worker.DB.Model(&models.PDDProductReplacementWorkerLease{}).Where("id = 1 AND owner_id = ?", worker.OwnerID).
Updates(map[string]any{"owner_id": "", "expires_at": time.Unix(0, 0).UTC()})
}
func (worker *Worker) defaults() {
if worker.Matcher == nil {
worker.Matcher = aimatching.NewService(worker.DB)
}
if worker.OwnerID == "" {
worker.OwnerID = uuid.NewString()
}
if worker.Now == nil {
worker.Now = func() time.Time { return time.Now().UTC() }
}
if worker.MaxAttempts < 1 {
worker.MaxAttempts = defaultMatchAttempts
}
if worker.LeaseDuration <= 0 {
worker.LeaseDuration = defaultWorkerLease
}
}
type pddCandidates struct {
colors, sizes []string
colorRoles, sizeRoles int
}
func replacementCandidates(raw string) (pddCandidates, error) {
var dimensions []struct {
Role string `json:"role"`
Values []struct {
Name string `json:"name"`
Selectable *bool `json:"selectable,omitempty"`
} `json:"values"`
}
if err := json.Unmarshal([]byte(raw), &dimensions); err != nil {
return pddCandidates{}, err
}
result := pddCandidates{}
for _, dimension := range dimensions {
if dimension.Role == shopeeproduct.RoleColor {
result.colorRoles++
}
if dimension.Role == shopeeproduct.RoleSize {
result.sizeRoles++
}
for _, value := range dimension.Values {
if value.Selectable != nil && !*value.Selectable {
continue
}
name := strings.TrimSpace(value.Name)
if name == "" {
continue
}
if dimension.Role == shopeeproduct.RoleColor {
result.colors = appendUnique(result.colors, name)
}
if dimension.Role == shopeeproduct.RoleSize {
result.sizes = appendUnique(result.sizes, name)
}
}
}
return result, nil
}
func uniqueMatchRoles(specs []shopeeproduct.SpecDimension) bool {
colorRoles, sizeRoles := 0, 0
for _, dimension := range specs {
if dimension.Role == shopeeproduct.RoleColor {
colorRoles++
}
if dimension.Role == shopeeproduct.RoleSize {
sizeRoles++
}
}
return colorRoles <= 1 && sizeRoles <= 1
}
func candidateValues(role string, candidates pddCandidates) []string {
if role == shopeeproduct.RoleColor {
return candidates.colors
}
return candidates.sizes
}
func containsCandidate(target string, candidates []string) bool {
for _, candidate := range candidates {
if target == candidate {
return true
}
}
return false
}
func cloneSpecs(specs []shopeeproduct.SpecDimension) []shopeeproduct.SpecDimension {
raw, _ := json.Marshal(specs)
cloned := []shopeeproduct.SpecDimension{}
_ = json.Unmarshal(raw, &cloned)
return cloned
}
func appendUnique(values []string, value string) []string {
for _, existing := range values {
if existing == value {
return values
}
}
return append(values, value)
}
func compactReason(value string, limit int) string {
value = strings.TrimSpace(value)
runes := []rune(value)
if len(runes) > limit {
value = string(runes[:limit])
}
return value
}
func matchErrorCode(err error) string {
var target *aimatching.Error
if errors.As(err, &target) {
return target.Code
}
return "MATCH_FAILED"
}
func retryableMatchError(err error) bool {
var target *aimatching.Error
return errors.As(err, &target) && target.Code == aimatching.CodeProviderUnavailable
}
@@ -0,0 +1,221 @@
package replacement
import (
"context"
"errors"
"testing"
"go-admin/app/goauto/aimatching"
"go-admin/app/goauto/models"
"go-admin/app/goauto/shopeeproduct"
"github.com/google/uuid"
"gorm.io/gorm"
)
type stubMatcher struct {
result aimatching.MatchResult
err error
calls int
}
func (matcher *stubMatcher) Resolve(_ context.Context, _ aimatching.MatchRequest) (aimatching.MatchResult, error) {
matcher.calls++
return matcher.result, matcher.err
}
func TestWorkerConfirmsExactMatchWithoutConfidence(t *testing.T) {
fixture, item := activatedMatchFixture(t)
matcher := &stubMatcher{result: aimatching.MatchResult{MappedColor: "红色", Source: aimatching.SourceExact}}
worker := NewWorker(fixture.db)
worker.Matcher = matcher
processed, err := worker.RunOnce(context.Background())
if err != nil || !processed || matcher.calls != 1 {
t.Fatalf("processed=%v calls=%d err=%v", processed, matcher.calls, err)
}
var got models.PDDProductReplacementItem
if err := fixture.db.First(&got, item.ID).Error; err != nil {
t.Fatal(err)
}
if got.MappingStatus != models.ReplacementItemMappingMatched || got.Source == nil || *got.Source != models.ReplacementItemSourceExactMatch || got.Confidence != nil {
t.Fatalf("item=%+v", got)
}
var product models.ShopeeProduct
if err := fixture.db.First(&product, item.ShopeeProductID).Error; err != nil {
t.Fatal(err)
}
specs, _ := shopeeproduct.Unmarshal(product.SpecsJSON)
if specs[0].Values[0].Mapping == nil || specs[0].Values[0].Mapping.Status != shopeeproduct.MappingStatusConfirmed {
t.Fatalf("specs=%+v", specs)
}
}
func TestWorkerRejectsAIMatchBelowThresholdAndDoesNotOverwriteChangedInput(t *testing.T) {
t.Run("below threshold", func(t *testing.T) {
fixture, item := activatedMatchFixture(t)
confidence := 0.89
matcher := &stubMatcher{result: aiResult("红色", &confidence, "候选中最接近")}
worker := NewWorker(fixture.db)
worker.Matcher = matcher
if _, err := worker.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
assertManualRequired(t, fixture.db, item.ID, "AI_AUTO_CONFIRM_REJECTED")
})
t.Run("input changed during match", func(t *testing.T) {
fixture, item := activatedMatchFixture(t)
matcher := &mutatingMatcher{db: fixture.db, shopeeID: item.ShopeeProductID}
worker := NewWorker(fixture.db)
worker.Matcher = matcher
if _, err := worker.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
assertManualRequired(t, fixture.db, item.ID, "MATCH_INPUT_CHANGED")
var product models.ShopeeProduct
_ = fixture.db.First(&product, item.ShopeeProductID).Error
if product.SpecsJSON != `[]` {
t.Fatalf("worker overwrote concurrent manual change: %s", product.SpecsJSON)
}
})
}
func TestWorkerRejectsInvalidAIAutoConfirmationEvidence(t *testing.T) {
below := 0.89
aboveRange := 1.01
tests := []struct {
name string
result aimatching.MatchResult
code string
}{
{name: "missing confidence", result: aiResult("红色", nil, "候选唯一"), code: "AI_AUTO_CONFIRM_REJECTED"},
{name: "out of range", result: aiResult("红色", &aboveRange, "候选唯一"), code: "AI_AUTO_CONFIRM_REJECTED"},
{name: "below threshold", result: aiResult("红色", &below, "候选唯一"), code: "AI_AUTO_CONFIRM_REJECTED"},
{name: "empty reason", result: aiResult("红色", floatPointer(0.95), ""), code: "AI_AUTO_CONFIRM_REJECTED"},
{name: "outside candidates", result: aiResult("相近红", floatPointer(0.95), "相近"), code: "MATCH_OUTSIDE_CANDIDATES"},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
fixture, item := activatedMatchFixture(t)
worker := NewWorker(fixture.db)
worker.Matcher = &stubMatcher{result: test.result}
if _, err := worker.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
assertManualRequired(t, fixture.db, item.ID, test.code)
})
}
}
func TestWorkerDoesNotCallMatcherWhenAIIsDisabled(t *testing.T) {
fixture, item := activatedMatchFixture(t)
if err := fixture.db.Model(&models.AIMatchingSetting{}).Where("id = 1").Update("enabled", false).Error; err != nil {
t.Fatal(err)
}
matcher := &stubMatcher{result: aimatching.MatchResult{MappedColor: "红色", Source: aimatching.SourceExact}}
worker := NewWorker(fixture.db)
worker.Matcher = matcher
if _, err := worker.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if matcher.calls != 0 {
t.Fatalf("matcher calls=%d", matcher.calls)
}
assertManualRequired(t, fixture.db, item.ID, aimatching.CodeNotConfigured)
}
func TestWorkerRetriesProviderFailureOnlyToFiniteLimit(t *testing.T) {
fixture, item := activatedMatchFixture(t)
matcher := &stubMatcher{err: &aimatching.Error{Code: aimatching.CodeProviderUnavailable, Message: "down", Cause: errors.New("temporary")}}
worker := NewWorker(fixture.db)
worker.Matcher = matcher
worker.MaxAttempts = 3
if err := worker.RunPending(context.Background()); err != nil {
t.Fatal(err)
}
if matcher.calls != 3 {
t.Fatalf("calls=%d", matcher.calls)
}
assertManualRequired(t, fixture.db, item.ID, aimatching.CodeProviderUnavailable)
}
func TestManualMappingCompletesReplacementItemAndAggregate(t *testing.T) {
fixture, item := activatedMatchFixture(t)
if err := fixture.db.Model(&models.PDDProductReplacementItem{}).Where("id = ?", item.ID).Update("mapping_status", models.ReplacementItemMappingManualRequired).Error; err != nil {
t.Fatal(err)
}
service := shopeeproduct.NewService(fixture.db)
if _, err := service.SetMapping(context.Background(), item.ShopeeProductID, shopeeproduct.SetMappingRequest{
RequestID: uuid.NewString(), Dimension: "颜色", ValueName: "红色", PDDValue: "红色", Source: shopeeproduct.MappingSourceManual,
}); err != nil {
t.Fatal(err)
}
var got models.PDDProductReplacementItem
if err := fixture.db.First(&got, item.ID).Error; err != nil {
t.Fatal(err)
}
if got.MappingStatus != models.ReplacementItemMappingMatched || got.Source == nil || *got.Source != models.ReplacementItemSourceManualMapping {
t.Fatalf("item=%+v", got)
}
var parent models.PDDProductReplacement
if err := fixture.db.First(&parent, item.ReplacementID).Error; err != nil || parent.MappingStatus != models.ReplacementMappingCompleted {
t.Fatalf("parent=%+v err=%v", parent, err)
}
}
type mutatingMatcher struct {
db *gorm.DB
shopeeID uint64
}
func (matcher *mutatingMatcher) Resolve(_ context.Context, _ aimatching.MatchRequest) (aimatching.MatchResult, error) {
if err := matcher.db.Model(&models.ShopeeProduct{}).Where("id = ?", matcher.shopeeID).Update("specs_json", `[]`).Error; err != nil {
return aimatching.MatchResult{}, err
}
return aimatching.MatchResult{MappedColor: "红色", Source: aimatching.SourceExact}, nil
}
func activatedMatchFixture(t *testing.T) (replacementFixture, models.PDDProductReplacementItem) {
t.Helper()
fixture := seedReplacementFixture(t)
fixture.service.StartMatching = func(*gorm.DB) {}
if err := fixture.db.Model(&models.PDDProduct{}).Where("id = ?", fixture.target.ID).Update("specs_json", `[{"name":"颜色","role":"color","values":[{"name":"红色"}]}]`).Error; err != nil {
t.Fatal(err)
}
shopee := models.ShopeeProduct{ShopeeItemID: "MATCH-1", PDDProductID: &fixture.source.ID, Currency: "CNY", SpecsJSON: `[{"name":"颜色","role":"color","values":[{"name":"红色","source":"import"}]}]`}
if err := fixture.db.Create(&shopee).Error; err != nil {
t.Fatal(err)
}
result, err := fixture.service.RegisterAndActivate(context.Background(), fixture.request())
if err != nil {
t.Fatal(err)
}
if len(result.Record.Items) != 1 {
t.Fatalf("items=%d", len(result.Record.Items))
}
setting := models.AIMatchingSetting{ID: 1, Enabled: true, Provider: aimatching.ProviderOpenAICompatible, BaseURL: "https://example.invalid/v1", Model: "test", APIKey: "test-only", TimeoutSeconds: 15, AutoConfirmMinConfidence: 0.9}
if err := fixture.db.Create(&setting).Error; err != nil {
t.Fatal(err)
}
return fixture, result.Record.Items[0]
}
func aiResult(color string, confidence *float64, reason string) aimatching.MatchResult {
result := aimatching.MatchResult{MappedColor: color, Source: aimatching.SourceAI}
result.Decision.Confidence = confidence
result.Decision.Reason = reason
return result
}
func floatPointer(value float64) *float64 { return &value }
func assertManualRequired(t *testing.T, db *gorm.DB, itemID uint64, code string) {
t.Helper()
var item models.PDDProductReplacementItem
if err := db.First(&item, itemID).Error; err != nil {
t.Fatal(err)
}
if item.MappingStatus != models.ReplacementItemMappingManualRequired || item.LastErrorCode == nil || *item.LastErrorCode != code {
t.Fatalf("item=%+v", item)
}
}
+96 -24
View File
@@ -5,6 +5,7 @@ import (
"errors"
"fmt"
"strings"
"time"
adminmodels "go-admin/app/admin/models"
"go-admin/app/goauto/models"
@@ -492,37 +493,108 @@ func (service *Service) mutateSpecs(ctx context.Context, id uint64, requestID st
return SaveResponse{}, internalError(err)
}
var record models.ShopeeProduct
if err := db.First(&record, id).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return SaveResponse{}, productNotFound()
err := db.Transaction(func(tx *gorm.DB) error {
var record models.ShopeeProduct
if err := tx.First(&record, id).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return productNotFound()
}
return internalError(err)
}
return SaveResponse{}, internalError(err)
}
specs, err := Unmarshal(record.SpecsJSON)
if err != nil {
return SaveResponse{}, internalError(err)
}
updated, err := mutate(specs)
specs, err := Unmarshal(record.SpecsJSON)
if err != nil {
return internalError(err)
}
updated, err := mutate(specs)
if err != nil {
return err
}
if err := Validate(updated); err != nil {
return invalidRequest(err.Error())
}
specsJSON, err := Marshal(updated)
if err != nil {
return internalError(err)
}
if err := tx.Model(&models.ShopeeProduct{}).Where("id = ?", id).Updates(map[string]any{
"specs_json": specsJSON, "last_update_request_id": requestID,
}).Error; err != nil {
return internalError(err)
}
if err := syncReplacementItemAfterManualEdit(tx, id, updated); err != nil {
return internalError(err)
}
return nil
})
if err != nil {
return SaveResponse{}, err
}
if err := Validate(updated); err != nil {
return SaveResponse{}, invalidRequest(err.Error())
}
specsJSON, err := Marshal(updated)
if err != nil {
return SaveResponse{}, internalError(err)
}
result := db.Model(&models.ShopeeProduct{}).Where("id = ?", id).Updates(map[string]any{
"specs_json": specsJSON, "last_update_request_id": requestID,
})
if result.Error != nil {
return SaveResponse{}, internalError(result.Error)
}
return service.Detail(ctx, id)
}
func syncReplacementItemAfterManualEdit(tx *gorm.DB, shopeeProductID uint64, specs []SpecDimension) error {
if !tx.Migrator().HasTable(&models.PDDProductReplacementItem{}) {
return nil
}
complete := true
for _, dimension := range specs {
if dimension.Role != RoleColor && dimension.Role != RoleSize {
continue
}
for _, value := range dimension.Values {
if value.Mapping == nil || value.Mapping.Status != MappingStatusConfirmed {
complete = false
}
}
}
var items []models.PDDProductReplacementItem
if err := tx.Joins("JOIN pdd_product_replacement r ON r.id = pdd_product_replacement_item.replacement_id").
Where("pdd_product_replacement_item.shopee_product_id = ? AND r.status = ?", shopeeProductID, models.ReplacementStatusActive).
Find(&items).Error; err != nil {
return err
}
for _, item := range items {
status := models.ReplacementItemMappingManualRequired
updates := map[string]any{"mapping_status": status, "updated_at": time.Now().UTC()}
if complete {
source := models.ReplacementItemSourceManualMapping
updates["mapping_status"], updates["source"], updates["confidence"], updates["last_error_code"], updates["last_error_at"] = models.ReplacementItemMappingMatched, source, nil, nil, nil
}
if err := tx.Model(&models.PDDProductReplacementItem{}).Where("id = ?", item.ID).Updates(updates).Error; err != nil {
return err
}
if err := syncReplacementAggregate(tx, item.ReplacementID); err != nil {
return err
}
}
return nil
}
func syncReplacementAggregate(tx *gorm.DB, replacementID uint64) error {
var items []models.PDDProductReplacementItem
if err := tx.Where("replacement_id = ?", replacementID).Find(&items).Error; err != nil {
return err
}
matched, terminal := 0, 0
for _, item := range items {
if item.MappingStatus == models.ReplacementItemMappingMatched {
matched++
}
if item.MappingStatus != models.ReplacementItemMappingMatching {
terminal++
}
}
status := models.ReplacementMappingMatching
if len(items) > 0 && terminal == len(items) {
status = models.ReplacementMappingCompletedPartial
if matched == len(items) {
status = models.ReplacementMappingCompleted
}
}
return tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PDDProductReplacement{}).Where("id = ?", replacementID).
Updates(map[string]any{"mapping_status": status, "mapping_updated_at": time.Now().UTC()}).Error
}
// ---------------------------------------------------------------- batch soft delete / restore
type BatchDeleteRequest struct {
+6 -6
View File
@@ -16,9 +16,9 @@ const (
ValueSourceManual = "manual"
)
// Mapping origin. Every mapping needs human confirmation before it takes
// effect, including exact name matches: an identical name does not guarantee an
// identical physical size (Taiwan versus mainland sizing, for example).
// Mapping origin. #131 allows deterministic exact matches and strictly gated
// high-confidence AI matches to be born confirmed; manual confirmation remains
// available for every other result.
const (
MappingSourceManual = "manual"
MappingSourceExactMatch = "exact_match"
@@ -129,10 +129,10 @@ func validateMapping(mapping *Mapping) error {
default:
return invalid("映射状态无效")
}
// Only manual mappings may be born confirmed. Exact match and AI match are
// suggestions and must pass through human confirmation (#40, #46).
if mapping.Status == MappingStatusConfirmed && mapping.Source == MappingSourceAIMatch {
return invalid("AI 匹配结果必须人工确认后才能置为已确认")
if mapping.Confidence == nil || strings.TrimSpace(mapping.Reason) == "" {
return invalid("自动确认的 AI 映射必须记录置信度和原因")
}
}
if mapping.Confidence != nil && (*mapping.Confidence < 0 || *mapping.Confidence > 1) {
return invalid("匹配置信度必须在 0 与 1 之间")
+2
View File
@@ -21,6 +21,7 @@ import (
"go-admin/app/admin/models"
"go-admin/app/admin/router"
goautodevice "go-admin/app/goauto/device"
goautoreplacement "go-admin/app/goauto/replacement"
goautosybimport "go-admin/app/goauto/sybimport"
goautosybinnercode "go-admin/app/goauto/sybinnercode"
"go-admin/app/jobs"
@@ -98,6 +99,7 @@ func run() error {
if err := goautosybinnercode.RecoverInterrupted(db); err != nil {
return fmt.Errorf("recover interrupted SYB inner-code writes: %w", err)
}
goautoreplacement.RecoverMatching(db)
}
offlineMonitorContext, stopOfflineMonitors := context.WithCancel(context.Background())
defer stopOfflineMonitors()
@@ -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), migratePDDReplacementActivation)
}
func migratePDDReplacementActivation(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
})
}