diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index 9be7041..3d59669 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Architecture-and-Code-Map wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.- -wiki_revision: d585a906857a46752788c9f25d3ac0361739dada -synchronized_at: 2026-08-28T09:27:21Z +wiki_revision: b54f4cce004cb6eb2398b40f89ad88888b773a3a +synchronized_at: 2026-08-28T09:50:04Z # 架构与代码地图 @@ -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。 diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 56009f5..5643564 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Business-Rules-and-Glossary wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.- -wiki_revision: b31ebf8888da7c3d9b71520938b104ff9e66511f -synchronized_at: 2026-08-28T09:27:33Z +wiki_revision: 14d111bf1043f9c31f4f54dbfbca38826e986ab9 +synchronized_at: 2026-08-28T09:50:14Z # 业务规则与术语 @@ -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,并只处理原记录冻结的虾皮商品影响集合。 diff --git a/server/app/goauto/aimatching/service.go b/server/app/goauto/aimatching/service.go index d96868b..4255494 100644 --- a/server/app/goauto/aimatching/service.go +++ b/server/app/goauto/aimatching/service.go @@ -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(¤t).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 { diff --git a/server/app/goauto/aimatching/service_test.go b/server/app/goauto/aimatching/service_test.go index 880ba0b..0238445 100644 --- a/server/app/goauto/aimatching/service_test.go +++ b/server/app/goauto/aimatching/service_test.go @@ -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) + } +} diff --git a/server/app/goauto/aimatching/types.go b/server/app/goauto/aimatching/types.go index d9a9447..0c10e31 100644 --- a/server/app/goauto/aimatching/types.go +++ b/server/app/goauto/aimatching/types.go @@ -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 { diff --git a/server/app/goauto/migrations/migrate.go b/server/app/goauto/migrations/migrate.go index 46ef634..3079e50 100644 --- a/server/app/goauto/migrations/migrate.go +++ b/server/app/goauto/migrations/migrate.go @@ -43,6 +43,7 @@ func MigratedModels() []any { &models.CollectionSKUValue{}, &models.PDDProductReplacement{}, &models.PDDProductReplacementItem{}, + &models.PDDProductReplacementWorkerLease{}, } } diff --git a/server/app/goauto/models/replacement.go b/server/app/goauto/models/replacement.go index 2b03334..45668dc 100644 --- a/server/app/goauto/models/replacement.go +++ b/server/app/goauto/models/replacement.go @@ -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" +} diff --git a/server/app/goauto/models/schema.go b/server/app/goauto/models/schema.go index 41d930f..dca7f3c 100644 --- a/server/app/goauto/models/schema.go +++ b/server/app/goauto/models/schema.go @@ -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" } diff --git a/server/app/goauto/replacement/activation.go b/server/app/goauto/replacement/activation.go new file mode 100644 index 0000000..159b7fd --- /dev/null +++ b/server/app/goauto/replacement/activation.go @@ -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 +} diff --git a/server/app/goauto/replacement/activation_test.go b/server/app/goauto/replacement/activation_test.go new file mode 100644 index 0000000..9a0c5c7 --- /dev/null +++ b/server/app/goauto/replacement/activation_test.go @@ -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 +} diff --git a/server/app/goauto/replacement/service.go b/server/app/goauto/replacement/service.go index 604fe11..b40fa71 100644 --- a/server/app/goauto/replacement/service.go +++ b/server/app/goauto/replacement/service.go @@ -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, "来源采购任务不满足替换条件") } } diff --git a/server/app/goauto/replacement/worker.go b/server/app/goauto/replacement/worker.go new file mode 100644 index 0000000..12f7939 --- /dev/null +++ b/server/app/goauto/replacement/worker.go @@ -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(¤t, 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(¤t, 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 +} diff --git a/server/app/goauto/replacement/worker_test.go b/server/app/goauto/replacement/worker_test.go new file mode 100644 index 0000000..7bb9a8d --- /dev/null +++ b/server/app/goauto/replacement/worker_test.go @@ -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) + } +} diff --git a/server/app/goauto/shopeeproduct/service.go b/server/app/goauto/shopeeproduct/service.go index c12f87d..5beada7 100644 --- a/server/app/goauto/shopeeproduct/service.go +++ b/server/app/goauto/shopeeproduct/service.go @@ -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 { diff --git a/server/app/goauto/shopeeproduct/specs.go b/server/app/goauto/shopeeproduct/specs.go index bb3d35a..77c6880 100644 --- a/server/app/goauto/shopeeproduct/specs.go +++ b/server/app/goauto/shopeeproduct/specs.go @@ -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 之间") diff --git a/server/cmd/api/server.go b/server/cmd/api/server.go index bd54c34..bb63d48 100644 --- a/server/cmd/api/server.go +++ b/server/cmd/api/server.go @@ -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() diff --git a/server/cmd/migrate/migration/version-local/1787885100000_pdd_replacement_activation.go b/server/cmd/migrate/migration/version-local/1787885100000_pdd_replacement_activation.go new file mode 100644 index 0000000..3cb27a3 --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1787885100000_pdd_replacement_activation.go @@ -0,0 +1,25 @@ +package version_local + +import ( + "runtime" + + goautomigrations "go-admin/app/goauto/migrations" + "go-admin/cmd/migrate/migration" + common "go-admin/common/models" + + "gorm.io/gorm" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), 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 + }) +}