merge: resume bounded Shopee auto-match scans (#359)
This commit is contained in:
@@ -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: 3165d7d419a46f41fa70a63286799deda9a8ce9f
|
||||
synchronized_at: 2026-10-06T09:50:17Z
|
||||
wiki_revision: b8b969b0c20aca9af6798a5a0b8784aa735600d4
|
||||
synchronized_at: 2026-10-07T08:07:14Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
<!-- gitea-wiki-mirror:start -->
|
||||
@@ -627,3 +627,15 @@ Web 唯一展示位置为“采集采购 → SYB 同步记录”:列表状态
|
||||
- PddProductDetailCollector 的入口、快速确认恢复和颜色点击通过 clickFreshDetailed 获取结果;SpecClickDiagnostic 构造白名单结构事件。AgentForegroundService 给探测执行器与采集器接入既有 SafeAgentDiagnosticRecorder / 单线程队列,未接入原始 trace。
|
||||
- AgentDiagnosticSchema v3:agent_diagnostic 追加可空 task_type TEXT、task_attempt_id TEXT、device_id INTEGER、phase TEXT、rule_snapshot_hash TEXT。onUpgrade 支持 V1/V2 追加并检查已有列;保留旧行且新字段为 NULL。AgentDiagnosticStore 写采购记录时验证类型、UUID、正设备 ID、阶段和 64 位十六进制哈希。动作 attempt 与采购 attempt UUID 分离,全库 50 条/7 天保留边界不变。
|
||||
- 无 Server/Web/业务库或共享接口字段变化,无订单提交流程变化。新错误沿既有 errorCode 字符串回传;真实探测/采购验收仍待用户授权。
|
||||
|
||||
## 自动匹配扫描游标与租约守卫(#359)
|
||||
|
||||
实现绑定 `9fcbc64117bcee0cbed25c3957a14a25d826f637`,仅开发分支完成,迁移与部署尚未执行。
|
||||
|
||||
- `server/app/goauto/shopeeproduct/auto_match_batch.go` 按商品ID键集分页,200/页、2000/轮、10分钟预算,复用原单商品匹配。30分钟运行租约与唯一active_slot不变,运行及工作项变更增加所有权检查。
|
||||
- `shopee_spec_auto_match_run.resume_after_id` 为可空、非负BIGINT:NULL不提交位置,0回绕,从最近已终结非NULL运行读取;`stop_reason` 为VARCHAR(24)、NOT NULL DEFAULT ''。完成更新在相同所有权守卫下原子提交统计和位置。
|
||||
- `1791300000000_shopee_spec_auto_match_resume.go` 只追加两列,重复执行幂等;既有运行初始化NULL/空字符串,不改商品或工作项,不改变定时任务启停。旧代码忽略新列,回退代码保留列和既有映射。
|
||||
- 批处理私有context将所有权检查传递到`ai_suggest.go`的Provider调用以及`auto_match.go`的映射事务;复用当前事务锁定运行,非批处理上下文不引入运行查询。AI决策算法和匹配规则不变。
|
||||
- 因预算超时不能再使用已取消context写统计,收尾仅使用最多5秒的独立上下文执行受所有权保护的完成更新,不启动新商品领取或AI调用;失租不强制落库。
|
||||
- MySQL默认返回实际修改行数;续租更新返回0时,只在当前持有行锁的事务内再次核验owner/状态/槽位/实时有效租约,以区分同毫秒值未变化与真实失租;其他完成/工作项更新仍要求恰好一行。
|
||||
- AI配置读取先返回数据库错误,再判断停用,避免基础设施错误被误记为业务跳过。仅批次私有上下文把基础设施错误作为本轮错误终止;普通单商品Provider重试策略保持不变。
|
||||
|
||||
@@ -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: d9b2a343e790c7d90354d47e1a977f134a72a760
|
||||
synchronized_at: 2026-10-07T03:28:45Z
|
||||
wiki_revision: 2abd7a07ccf8a28308ce9c318250e2f8f957c137
|
||||
synchronized_at: 2026-10-07T08:07:17Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
<!-- gitea-wiki-mirror:start -->
|
||||
@@ -824,3 +824,14 @@ Android 0.9.64 / versionCode 77,源码 `6550b9f`(分支实现,尚未安装
|
||||
- 合并后仍使用外层完整原文按既有规则去尾价作为规格值,沿用原点击节点排序、selected/checked 聚合;普通采集与采购探测共用这一路径。精确定位、即时确认、选中证明、最终复核及下单行为不改,不允许模糊点击。
|
||||
- 历史任务若已映射到截短值,不自动升格为完整值,不修改任务快照、映射或历史数据;找不到精确目标仍明确失败。即时确认先查其他选中值,而最终确认先接受唯一目标已选中,是既有实现差异,本修复不调整或掩盖该差异。
|
||||
- #361 的诊断未写入疑点继续独立核查。解析合成测试可先行,后续用于手机安装的集成版本须包含 #361,并经明确授权安装/真机验证;不以已有单次采购成功代替完整验收,不执行付款。
|
||||
|
||||
## 蝦皮规格自动匹配有界续扫(#359)
|
||||
|
||||
实现绑定 `9fcbc64117bcee0cbed25c3957a14a25d826f637`,分支 `fix/359-auto-match-resume`;已实现并经本地合成验证,尚未合并、执行业务库迁移或发布,不代表当前线上已启用。
|
||||
|
||||
- 定时及管理员批量匹配沿用原匹配算法、阈值、人工/有效确认映射保护和指纹重试规则;仅修复固定首段扫描无法到达后方候选。单商品手动匹配不受批次租约检查影响。
|
||||
- 使用商品 ID 升序键集分页,每页最多200件,每轮实际检查最多2000件,默认实际领取处理最多20件(batchLimit原校验范围不变)。整轮数据库及AI操作共享10分钟预算,逐商品串行,30分钟租约不变。
|
||||
- SQL排除明确空规格,Go先检查蝦皮端非空颜色/尺码再读PDD;只有颜色或只有尺码仍合法,只有other/空values不能成为匹配候选。
|
||||
- 正常完成或预算退出只保存最后已确定处理/跳过的位置;页中提前退出不跳到预取末尾。确实消费完末页才回绕0;下轮/进程重启从最近已终结且有有效游标的运行续扫,NULL不是有效游标,0是有效回绕点。
|
||||
- 单运行所有权在分页续期、领取、Provider调用及保存映射时检查。失租旧运行不能继续领取或覆盖新owner;基础设施错误或失租不提交新游标。
|
||||
- completed只表示本轮正常结束,处理0件可能合法;scanned为实际检查数,不是预取数或全表数,processed不是成功数,confirmed/unmatched是规格项数。商品变化后可能需要等扫描回绕,不保证固定小时内全部处理。
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Android-Agent-API-Contract
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Android-Agent-API-Contract.-
|
||||
wiki_revision: 1ad78ac33cf5876c6b7db6a89107854978066373
|
||||
synchronized_at: 2026-10-06T01:39:51Z
|
||||
wiki_revision: cf3788145daf3dfe9a4526fd231bc4b8d730f763
|
||||
synchronized_at: 2026-10-07T08:07:41Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
<!-- gitea-wiki-mirror:start -->
|
||||
@@ -1472,3 +1472,13 @@ Web 输入去重后一个值提交旧标量、多个值提交重复集合键,
|
||||
| sybStatusSyncedAt | RFC3339 时间字符串或 null | 最近一次正常同步取得有效取消值的 UTC 时间,不代表上游取消发生时间 |
|
||||
|
||||
只使用上游 isCancel;不兼容性猜测字符串/布尔/其他数值,缺失或无效值保留原字段。成功同步的 false 可覆盖 true。无有效值的新记录两个字段均为 null;旧客户端可忽略新增字段,新 Web 对旧响应缺字段显示未获取。列表与详情不触发额外上游请求。取消状态仅展示,不改变采购准备阶段、创建/重试资格、任务执行或现有订单事实。
|
||||
|
||||
## Admin 蝦皮规格自动匹配运行摘要追加字段(#359)
|
||||
|
||||
实现绑定 `9fcbc64117bcee0cbed25c3957a14a25d826f637`,分支实现尚未部署。本节仅扩展既有Admin批次接口,不修改Android Agent接口、Web页面或权限。
|
||||
|
||||
- `POST /api/admin/v1/shopee-spec-auto-match/runs` 与 `GET /api/admin/v1/shopee-spec-auto-match/runs/latest` 的既有运行对象增加`resumeAfterId`、`stopReason`,请求参数、原字段和状态保持兼容。
|
||||
- `resumeAfterId`:可空非负整数;null表示该运行没有提交有效续扫点,0表示下一轮从头扫描,正数表示最后已完成检查的位置,不是预取页末商品。运行中/旧记录可能为null。
|
||||
- `stopReason`:旧记录默认空字符串;完成原因是`batch_limit`、`scan_budget`、`time_budget`、`end_of_scan`、`lease_lost`或`error`。失租旧进程不能为填此字段越权更新;合法回收路径标记lease_lost。
|
||||
- `scannedCount`改为实际检查的候选数;SQL已过滤的空档案及预取未检查项不计入。eligibleCount为Go资格通过数,processedCount为实际领取处理数;confirmedCount/unmatchedCount仍为规格项数,不能据此直接混算商品成功率。
|
||||
- 单轮处理默认20、每页200、实际扫描上限2000、整轮预算10分钟。正常0处理仍可completed,预算退出有持久游标;错误/失租不提交新游标。
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Deployment-and-Operations
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Deployment-and-Operations.-
|
||||
wiki_revision: 21e8ed0578212b84a77af537e13be6d10c8b5463
|
||||
synchronized_at: 2026-10-07T03:28:51Z
|
||||
wiki_revision: cf7b9d1e8ca8f97933936e4ddfd9edc3fac43c98
|
||||
synchronized_at: 2026-10-07T08:07:26Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
<!-- gitea-wiki-mirror:start -->
|
||||
@@ -284,3 +284,14 @@ Provider 故障日志只允许记录调用关联 ID、操作类型、耗时、
|
||||
- 公网首页、index.html、login、档口入库码路由均与新 dist/index.html 字节一致;10 项入口 JS/CSS 与健康接口正常,日志检查无 panic/fatal/1146/1054。真实 Chrome 已认证只读检查:首次默认 page=1/pageSize=200、当天空态与200条/页正常;清除日期后200行渲染通过,查询至渲染约987ms。未发起匹配、删除或回写,未保存原始生产数据或截图。
|
||||
- Web包 SHA256 `a6b4429280fae37e46e81bd21a594941cfee3b5e8c48c105a09254314af32311`;Server SHA256 `a4cf6c7028a2bbbfd86a1d657ad9c3daad7b4d08445405d8ca9a778adc584760`(未变)。
|
||||
- 回滚目录 `/home/goauto/releases/20261006-e76de6f-358` 保留。对本次纯 Web 变更可将 current 原子切回,不需要恢复数据库、删除文件或中断任务;若此后已升级后端,不能复用这一免重启结论。
|
||||
|
||||
## #359 续扫版本迁移与验证边界
|
||||
|
||||
源码 `9fcbc64117bcee0cbed25c3957a14a25d826f637`(`fix/359-auto-match-resume`)已实现;尚未执行本地/线上业务库迁移、合并、发布或真实匹配。以下为获得相应授权后的步骤,不是已上线结论。
|
||||
|
||||
1. 复核无冲突运行及现有迁移版本,按既有受限备份流程备份。追加迁移`1791300000000_shopee_spec_auto_match_resume.go`仅新增运行游标与停止原因,必须先迁移再运行新版本;不修改定时任务配置或历史商品。
|
||||
2. 发布后观察运行的stopReason/resumeAfterId及真实计数,确认多轮向后推进、末尾回绕,而非反复固定首段。合法无候选仍允许processedCount=0,不能要求每轮强制匹配成功。
|
||||
3. 批次结构化日志按run_id记录停止原因、实际扫描/领取/规格项计数及跳过类别,不记录商品规格原文或Provider响应。lost lease旧进程不能覆盖新运行,合法过期回收标记failed/lease_lost且不提交游标。
|
||||
4. 回退旧二进制时保留追加列和已保存映射,不通过数据库回滚覆盖后续业务。恢复处理会带来原本预期的AI调用和映射写入,仍受默认20件、串行与总时间预算限制。
|
||||
|
||||
本地测试使用SQLite内存库和模拟Provider;MySQL8.4.3只执行合成JSON粗过滤SELECT验证,未执行MySQL迁移/锁竞争集成验证,也未调用线上AI。MySQL正式迁移及真实批次验证仍须授权。
|
||||
|
||||
@@ -249,12 +249,15 @@ func (s *Service) ResolveSYBSpec(ctx context.Context, request SYBSpecParseReques
|
||||
|
||||
func (s *Service) activeSetting(ctx context.Context) (models.AIMatchingSetting, string, error) {
|
||||
setting, err := s.setting(ctx)
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) || !setting.Enabled {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return models.AIMatchingSetting{}, "", fail(CodeNotConfigured, "AI 规格匹配未启用")
|
||||
}
|
||||
if err != nil {
|
||||
return models.AIMatchingSetting{}, "", err
|
||||
}
|
||||
if !setting.Enabled {
|
||||
return models.AIMatchingSetting{}, "", fail(CodeNotConfigured, "AI 规格匹配未启用")
|
||||
}
|
||||
if strings.TrimSpace(setting.APIKey) == "" {
|
||||
return models.AIMatchingSetting{}, "", fail(CodeNotConfigured, "AI 规格匹配未配置 API Key")
|
||||
}
|
||||
|
||||
@@ -26,6 +26,9 @@ type ShopeeSpecAutoMatchRun struct {
|
||||
FinishedAt *time.Time `json:"finishedAt,omitempty"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
// NULL means no committed checkpoint; zero explicitly restarts at the head.
|
||||
ResumeAfterID *uint64 `json:"resumeAfterId" gorm:"type:bigint;check:ck_shopee_spec_auto_match_resume,resume_after_id IS NULL OR resume_after_id >= 0"`
|
||||
StopReason string `json:"stopReason" gorm:"size:24;not null;default:''"`
|
||||
}
|
||||
|
||||
func (ShopeeSpecAutoMatchRun) TableName() string { return "shopee_spec_auto_match_run" }
|
||||
|
||||
@@ -194,10 +194,33 @@ func (service *Service) suggestMappings(ctx context.Context, id uint64, requestC
|
||||
Dimension: role, ShopeeTitle: shopee.Title, PDDTitle: pdd.Title,
|
||||
Sources: sources, Candidates: candidates,
|
||||
}
|
||||
result, err := aiService.SuggestBatch(ctx, suggestReq)
|
||||
// Batch ownership can change between role calls or provider retries.
|
||||
var guardErr error
|
||||
suggest := func(request aimatching.SuggestRequest) (aimatching.SuggestResult, error) {
|
||||
if err := checkAutoMatchRunContext(ctx, service.DB); err != nil {
|
||||
guardErr = err
|
||||
return aimatching.SuggestResult{}, err
|
||||
}
|
||||
result, err := aiService.SuggestBatch(ctx, request)
|
||||
if _, batch := ctx.Value(autoMatchRunContextKey{}).(autoMatchRunGuard); batch && err != nil {
|
||||
var providerErr *aimatching.Error
|
||||
if !errors.As(err, &providerErr) {
|
||||
guardErr = internalError(err)
|
||||
return aimatching.SuggestResult{}, guardErr
|
||||
}
|
||||
}
|
||||
return result, err
|
||||
}
|
||||
result, err := suggest(suggestReq)
|
||||
if guardErr != nil {
|
||||
return AISuggestResponse{}, guardErr
|
||||
}
|
||||
suggestCalls := 1
|
||||
if err != nil {
|
||||
result, err = aiService.SuggestBatch(ctx, suggestReq)
|
||||
result, err = suggest(suggestReq)
|
||||
if guardErr != nil {
|
||||
return AISuggestResponse{}, guardErr
|
||||
}
|
||||
suggestCalls++
|
||||
if err != nil {
|
||||
return AISuggestResponse{}, aiUnavailable(aiSuggestErrorMessage(err))
|
||||
@@ -216,10 +239,13 @@ func (service *Service) suggestMappings(ctx context.Context, id uint64, requestC
|
||||
}
|
||||
}
|
||||
if len(retrySources) > 0 && suggestCalls < 2 {
|
||||
retryResult, retryErr := aiService.SuggestBatch(ctx, aimatching.SuggestRequest{
|
||||
retryResult, retryErr := suggest(aimatching.SuggestRequest{
|
||||
Dimension: role, ShopeeTitle: shopee.Title, PDDTitle: pdd.Title,
|
||||
Sources: retrySources, Candidates: candidates,
|
||||
})
|
||||
if guardErr != nil {
|
||||
return AISuggestResponse{}, guardErr
|
||||
}
|
||||
if retryErr == nil {
|
||||
for _, source := range retrySources {
|
||||
if decision, ok := retryResult.Decisions[source.ID]; ok {
|
||||
|
||||
@@ -148,6 +148,9 @@ func (service *Service) autoMatchMappings(ctx context.Context, id uint64, reques
|
||||
|
||||
replayed := false
|
||||
err = db.Transaction(func(tx *gorm.DB) error {
|
||||
if err := checkAutoMatchRunContext(ctx, tx.Clauses(clause.Locking{Strength: "UPDATE"})); err != nil {
|
||||
return err
|
||||
}
|
||||
var current models.ShopeeProduct
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(¤t, id).Error; err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
|
||||
@@ -13,8 +13,10 @@ import (
|
||||
"go-admin/app/goauto/aimatching"
|
||||
"go-admin/app/goauto/models"
|
||||
|
||||
log "github.com/go-admin-team/go-admin-core/logger"
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -23,6 +25,10 @@ const (
|
||||
autoMatchLeaseDuration = 30 * time.Minute
|
||||
autoMatchRetryDelay = time.Hour
|
||||
maxAutoMatchAttempts = 3
|
||||
autoMatchPageSize = 200
|
||||
autoMatchScanBudget = 2000
|
||||
autoMatchTimeBudget = 10 * time.Minute
|
||||
autoMatchNonEmptySpecsSQL = "TRIM(CAST(shopee_product.specs_json AS CHAR)) <> ? AND TRIM(CAST(shopee_product.specs_json AS CHAR)) <> ? AND TRIM(CAST(shopee_product.specs_json AS CHAR)) <> ? AND TRIM(CAST(shopee_product.specs_json AS CHAR)) <> ?"
|
||||
)
|
||||
|
||||
type AutoMatchRunView struct {
|
||||
@@ -41,11 +47,11 @@ func (service *Service) StartAutoMatchRun(ctx context.Context, trigger, requestI
|
||||
if trigger != "manual" && trigger != "scheduled" {
|
||||
return AutoMatchRunView{}, false, invalidRequest("trigger 无效")
|
||||
}
|
||||
if batchLimit <= 0 {
|
||||
if batchLimit == 0 {
|
||||
batchLimit = defaultAutoMatchBatchLimit
|
||||
}
|
||||
if batchLimit > 100 {
|
||||
return AutoMatchRunView{}, false, invalidRequest("batchLimit 不能超过 100")
|
||||
if batchLimit < 1 || batchLimit > 100 {
|
||||
return AutoMatchRunView{}, false, invalidRequest("batchLimit 必须在 1 到 100 之间")
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
lease := now.Add(autoMatchLeaseDuration)
|
||||
@@ -55,8 +61,8 @@ func (service *Service) StartAutoMatchRun(ctx context.Context, trigger, requestI
|
||||
created := false
|
||||
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Model(&models.ShopeeSpecAutoMatchRun{}).
|
||||
Where("status = ? AND active_slot = ? AND lease_expires_at < ?", "running", 1, now).
|
||||
Updates(map[string]any{"status": "failed", "active_slot": nil, "lease_owner": "", "lease_expires_at": nil, "error_summary": "上次运行租约过期,已安全释放", "finished_at": now}).Error; err != nil {
|
||||
Where("status = ? AND active_slot = ? AND lease_expires_at <= ?", "running", 1, now).
|
||||
Updates(map[string]any{"status": "failed", "active_slot": nil, "lease_owner": "", "lease_expires_at": nil, "error_summary": "上次运行租约过期,已安全释放", "finished_at": now, "resume_after_id": nil, "stop_reason": "lease_lost"}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if err := tx.Where("request_id = ?", requestID).First(&result).Error; err == nil {
|
||||
@@ -109,6 +115,8 @@ func (service *Service) LatestAutoMatchRun(ctx context.Context) (*AutoMatchRunVi
|
||||
// HTTP-launched goroutine or the scheduler because only the run owning the
|
||||
// active slot may update and finish itself.
|
||||
func (service *Service) ProcessAutoMatchRun(ctx context.Context, runID uint64) error {
|
||||
ctx, cancel := context.WithTimeout(ctx, autoMatchTimeBudget)
|
||||
defer cancel()
|
||||
var run models.ShopeeSpecAutoMatchRun
|
||||
if err := service.DB.WithContext(ctx).First(&run, runID).Error; err != nil {
|
||||
return err
|
||||
@@ -120,98 +128,149 @@ func (service *Service) ProcessAutoMatchRun(ctx context.Context, runID uint64) e
|
||||
if limit <= 0 || limit > 100 {
|
||||
limit = defaultAutoMatchBatchLimit
|
||||
}
|
||||
var candidates []models.ShopeeProduct
|
||||
queryLimit := limit * 25
|
||||
if queryLimit < 100 {
|
||||
queryLimit = 100
|
||||
ctx, cancelLease := context.WithCancelCause(ctx)
|
||||
defer cancelLease(nil)
|
||||
ctx = context.WithValue(ctx, autoMatchRunContextKey{}, autoMatchRunGuard{run: run, cancel: func() { cancelLease(errAutoMatchLeaseLost) }})
|
||||
stats := autoMatchBatchStats{}
|
||||
cursor := uint64(0)
|
||||
var previous models.ShopeeSpecAutoMatchRun
|
||||
err := service.DB.WithContext(ctx).Where("status <> ? AND resume_after_id IS NOT NULL", "running").Order("id DESC").First(&previous).Error
|
||||
if err == nil {
|
||||
cursor = *previous.ResumeAfterID
|
||||
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
if queryLimit > 1000 {
|
||||
queryLimit = 1000
|
||||
}
|
||||
if err := service.DB.WithContext(ctx).
|
||||
Joins("JOIN pdd_product ON pdd_product.id = shopee_product.pdd_product_id AND pdd_product.status = ?", "active").
|
||||
Where("shopee_product.pdd_product_id IS NOT NULL").
|
||||
Order("shopee_product.updated_at ASC, shopee_product.id ASC").Limit(queryLimit).Find(&candidates).Error; err != nil {
|
||||
service.finishAutoMatchRun(run, "failed", 0, 0, 0, 0, 0, 1, "扫描符合条件的商品失败")
|
||||
return err
|
||||
}
|
||||
|
||||
eligible, processed, confirmed, unmatched, failed := 0, 0, 0, 0, 0
|
||||
firstError := ""
|
||||
for _, product := range candidates {
|
||||
if processed >= limit {
|
||||
break
|
||||
stats.checkpointLoaded = true
|
||||
for {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "time_budget", err)
|
||||
}
|
||||
fingerprint, ok, err := service.autoMatchEligibility(ctx, product)
|
||||
if err != nil {
|
||||
failed++
|
||||
if firstError == "" {
|
||||
firstError = safeBatchError(err)
|
||||
if stats.processed >= limit {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "batch_limit", nil)
|
||||
}
|
||||
if stats.scanned >= autoMatchScanBudget {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "scan_budget", nil)
|
||||
}
|
||||
if err := service.renewAutoMatchRun(ctx, run); err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
pageLimit := min(autoMatchPageSize, autoMatchScanBudget-stats.scanned)
|
||||
var candidates []models.ShopeeProduct
|
||||
if err := service.DB.WithContext(ctx).
|
||||
Joins("JOIN pdd_product ON pdd_product.id = shopee_product.pdd_product_id AND pdd_product.status = ?", "active").
|
||||
Where("shopee_product.pdd_product_id IS NOT NULL AND shopee_product.id > ?", cursor).
|
||||
// Cast the JSON column to text before comparing: no JSON NOT IN/coercion.
|
||||
Where(autoMatchNonEmptySpecsSQL, "", "[]", "null", `""`).
|
||||
Order("shopee_product.id ASC").Limit(pageLimit).Find(&candidates).Error; err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
for _, product := range candidates {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "time_budget", err)
|
||||
}
|
||||
continue
|
||||
}
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
eligible++
|
||||
work, claimed, err := service.claimAutoMatchWork(ctx, run, product.ID, fingerprint)
|
||||
if err != nil {
|
||||
failed++
|
||||
if firstError == "" {
|
||||
firstError = safeBatchError(err)
|
||||
if stats.processed >= limit {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "batch_limit", nil)
|
||||
}
|
||||
continue
|
||||
}
|
||||
if !claimed {
|
||||
continue
|
||||
}
|
||||
processed++
|
||||
service.renewAutoMatchRun(run)
|
||||
response, matchErr := service.autoMatchMappings(ctx, product.ID, AutoMatchRequest{RequestID: uuid.NewString(), SpecContextVersion: fingerprint[:64]}, aimatching.MaxProviderTimeout)
|
||||
// fingerprint begins with the 64-character context version.
|
||||
postFingerprint := fingerprint
|
||||
if next, _, nextErr := service.autoMatchEligibility(ctx, product); nextErr == nil && next != "" {
|
||||
postFingerprint = next
|
||||
}
|
||||
if matchErr != nil {
|
||||
failed++
|
||||
if firstError == "" {
|
||||
firstError = safeBatchError(matchErr)
|
||||
stats.scanned++
|
||||
fingerprint, skip, err := service.autoMatchEligibilityReason(ctx, product)
|
||||
if err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
service.completeAutoMatchWork(work, postFingerprint, 0, 0, matchErr)
|
||||
continue
|
||||
if skip != "" {
|
||||
stats.skip(skip)
|
||||
cursor = product.ID
|
||||
continue
|
||||
}
|
||||
stats.eligible++
|
||||
work, claimed, err := service.claimAutoMatchWork(ctx, run, product.ID, fingerprint)
|
||||
if err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
if !claimed {
|
||||
stats.skip(autoMatchWorkSkip(work, fingerprint, time.Now().UTC()))
|
||||
cursor = product.ID
|
||||
continue
|
||||
}
|
||||
stats.processed++
|
||||
if err := checkAutoMatchRunContext(ctx, service.DB); err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
response, matchErr := service.autoMatchMappings(ctx, product.ID, AutoMatchRequest{RequestID: uuid.NewString(), SpecContextVersion: fingerprint[:64]}, aimatching.MaxProviderTimeout)
|
||||
if err := checkAutoMatchRunContext(ctx, service.DB); err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
if matchErr != nil && batchErrorCode(matchErr) == CodeInternal {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", matchErr)
|
||||
}
|
||||
postFingerprint := fingerprint
|
||||
if next, _, err := service.autoMatchEligibility(ctx, product); err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
} else if next != "" {
|
||||
postFingerprint = next
|
||||
}
|
||||
if err := service.completeAutoMatchWork(ctx, work, postFingerprint, response.ConfirmedCount, response.UnmatchedCount, matchErr); err != nil {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, cursor, "error", err)
|
||||
}
|
||||
if matchErr != nil {
|
||||
stats.failed++
|
||||
if stats.summary == "" {
|
||||
stats.summary = safeBatchError(matchErr)
|
||||
}
|
||||
} else {
|
||||
stats.confirmed += response.ConfirmedCount
|
||||
stats.unmatched += response.UnmatchedCount
|
||||
}
|
||||
cursor = product.ID
|
||||
}
|
||||
// Only a fully consumed short page proves the actual end of the scan.
|
||||
if len(candidates) < pageLimit {
|
||||
return service.endAutoMatchBatch(ctx, run, stats, 0, "end_of_scan", nil)
|
||||
}
|
||||
confirmed += response.ConfirmedCount
|
||||
unmatched += response.UnmatchedCount
|
||||
service.completeAutoMatchWork(work, postFingerprint, response.ConfirmedCount, response.UnmatchedCount, nil)
|
||||
}
|
||||
status := "completed"
|
||||
if failed > 0 {
|
||||
status = "completed_partial"
|
||||
}
|
||||
return service.finishAutoMatchRun(run, status, len(candidates), eligible, processed, confirmed, unmatched, failed, firstError)
|
||||
}
|
||||
|
||||
func (service *Service) autoMatchEligibility(ctx context.Context, product models.ShopeeProduct) (string, bool, error) {
|
||||
fingerprint, skip, err := service.autoMatchEligibilityReason(ctx, product)
|
||||
return fingerprint, skip == "" && err == nil, err
|
||||
}
|
||||
|
||||
func (service *Service) autoMatchEligibilityReason(ctx context.Context, product models.ShopeeProduct) (string, string, error) {
|
||||
if product.PDDProductID == nil {
|
||||
return "", false, nil
|
||||
return "", "no_specs", nil
|
||||
}
|
||||
var pdd models.PDDProduct
|
||||
if err := service.DB.WithContext(ctx).First(&pdd, *product.PDDProductID).Error; err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
if pdd.Status != "active" {
|
||||
return "", false, nil
|
||||
if strings.TrimSpace(product.SpecsJSON) == `""` {
|
||||
return "", "no_specs", nil
|
||||
}
|
||||
shopeeSpecs, err := Unmarshal(product.SpecsJSON)
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
return "", "", err
|
||||
}
|
||||
usable := false
|
||||
for _, dimension := range shopeeSpecs {
|
||||
if dimension.Role != RoleColor && dimension.Role != RoleSize {
|
||||
continue
|
||||
}
|
||||
for _, value := range dimension.Values {
|
||||
if strings.TrimSpace(value.Name) != "" {
|
||||
usable = true
|
||||
}
|
||||
}
|
||||
}
|
||||
if !usable {
|
||||
return "", "no_specs", nil
|
||||
}
|
||||
var pdd models.PDDProduct
|
||||
if err := service.DB.WithContext(ctx).First(&pdd, *product.PDDProductID).Error; err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
if pdd.Status != "active" {
|
||||
return "", "no_specs", nil
|
||||
}
|
||||
shared, needsMatch := false, false
|
||||
for _, role := range []string{RoleColor, RoleSize} {
|
||||
pddValues, err := selectablePDDValues(pdd.SpecsJSON, role)
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
return "", "", err
|
||||
}
|
||||
if len(pddValues) == 0 {
|
||||
continue
|
||||
@@ -220,24 +279,32 @@ func (service *Service) autoMatchEligibility(ctx context.Context, product models
|
||||
if dimension.Role != role || len(dimension.Values) == 0 {
|
||||
continue
|
||||
}
|
||||
shared = true
|
||||
for _, value := range dimension.Values {
|
||||
if strings.TrimSpace(value.Name) == "" {
|
||||
continue
|
||||
}
|
||||
shared = true
|
||||
if value.Mapping == nil || value.Mapping.Status != MappingStatusConfirmed || !pddValues[value.Mapping.PDDValue] {
|
||||
needsMatch = true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !shared || !needsMatch {
|
||||
return "", false, nil
|
||||
if !shared {
|
||||
return "", "no_specs", nil
|
||||
}
|
||||
if !needsMatch {
|
||||
return "", "confirmed", nil
|
||||
}
|
||||
contextVersion := computeSpecContextVersion(product.PDDProductID, product.SpecsJSON, pdd.SpecsJSON)
|
||||
var setting struct{ UpdatedAt time.Time }
|
||||
_ = service.DB.WithContext(ctx).Table((models.AIMatchingSetting{}).TableName()).Select("updated_at").Where("id = ?", 1).Scan(&setting).Error
|
||||
if err := service.DB.WithContext(ctx).Table((models.AIMatchingSetting{}).TableName()).Select("updated_at").Where("id = ?", 1).Scan(&setting).Error; err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
h := sha256.Sum256([]byte(contextVersion + "\x00" + setting.UpdatedAt.UTC().Format(time.RFC3339Nano)))
|
||||
// Keeping the context version as a prefix lets ProcessAutoMatchRun pass the
|
||||
// exact version to #194 without re-reading a potentially drifting input.
|
||||
return contextVersion + hex.EncodeToString(h[:]), true, nil
|
||||
return contextVersion + hex.EncodeToString(h[:]), "", nil
|
||||
}
|
||||
|
||||
func (service *Service) claimAutoMatchWork(ctx context.Context, run models.ShopeeSpecAutoMatchRun, productID uint64, fingerprint string) (models.ShopeeSpecAutoMatchWorkItem, bool, error) {
|
||||
@@ -245,6 +312,9 @@ func (service *Service) claimAutoMatchWork(ctx context.Context, run models.Shope
|
||||
lease := now.Add(autoMatchLeaseDuration)
|
||||
var work models.ShopeeSpecAutoMatchWorkItem
|
||||
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := lockAutoMatchRun(ctx, tx, run); err != nil {
|
||||
return err
|
||||
}
|
||||
err := tx.Where("shopee_product_id = ?", productID).First(&work).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
work = models.ShopeeSpecAutoMatchWorkItem{ShopeeProductID: productID, RunID: &run.ID, InputFingerprint: fingerprint, Status: "running", AttemptCount: 1, LeaseOwner: run.LeaseOwner, LeaseExpiresAt: &lease}
|
||||
@@ -253,11 +323,10 @@ func (service *Service) claimAutoMatchWork(ctx context.Context, run models.Shope
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if work.InputFingerprint == fingerprint {
|
||||
if work.Status == "completed" || work.Status == "unmatched" || work.AttemptCount >= maxAutoMatchAttempts || (work.NextAttemptAt != nil && work.NextAttemptAt.After(now)) || (work.Status == "running" && work.LeaseExpiresAt != nil && work.LeaseExpiresAt.After(now)) {
|
||||
return errWorkNotClaimed
|
||||
}
|
||||
} else {
|
||||
if autoMatchWorkSkip(work, fingerprint, now) != "" {
|
||||
return errWorkNotClaimed
|
||||
}
|
||||
if work.InputFingerprint != fingerprint {
|
||||
work.AttemptCount = 0
|
||||
}
|
||||
updates := map[string]any{"run_id": run.ID, "input_fingerprint": fingerprint, "status": "running", "attempt_count": work.AttemptCount + 1, "next_attempt_at": nil, "lease_owner": run.LeaseOwner, "lease_expires_at": lease, "last_error_code": "", "last_error": ""}
|
||||
@@ -274,7 +343,7 @@ func (service *Service) claimAutoMatchWork(ctx context.Context, run models.Shope
|
||||
|
||||
var errWorkNotClaimed = errors.New("auto match work not claimed")
|
||||
|
||||
func (service *Service) completeAutoMatchWork(work models.ShopeeSpecAutoMatchWorkItem, fingerprint string, confirmed, unmatched int, matchErr error) {
|
||||
func (service *Service) completeAutoMatchWork(ctx context.Context, work models.ShopeeSpecAutoMatchWorkItem, fingerprint string, confirmed, unmatched int, matchErr error) error {
|
||||
now := time.Now().UTC()
|
||||
updates := map[string]any{"input_fingerprint": fingerprint, "lease_owner": "", "lease_expires_at": nil, "confirmed_count": confirmed, "unmatched_count": unmatched}
|
||||
if matchErr == nil {
|
||||
@@ -294,18 +363,189 @@ func (service *Service) completeAutoMatchWork(work models.ShopeeSpecAutoMatchWor
|
||||
updates["next_attempt_at"] = nil
|
||||
}
|
||||
}
|
||||
_ = service.DB.Model(&models.ShopeeSpecAutoMatchWorkItem{}).Where("id = ?", work.ID).Updates(updates).Error
|
||||
return service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := checkAutoMatchRunContext(ctx, tx.Clauses(clause.Locking{Strength: "UPDATE"})); err != nil {
|
||||
return err
|
||||
}
|
||||
result := tx.Model(&models.ShopeeSpecAutoMatchWorkItem{}).Where("id = ? AND status = ? AND lease_owner = ? AND lease_expires_at > ?", work.ID, "running", work.LeaseOwner, now).Updates(updates)
|
||||
return autoMatchOwnedUpdate(result)
|
||||
})
|
||||
}
|
||||
|
||||
func (service *Service) renewAutoMatchRun(run models.ShopeeSpecAutoMatchRun) {
|
||||
lease := time.Now().UTC().Add(autoMatchLeaseDuration)
|
||||
_ = service.DB.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ? AND status = ? AND lease_owner = ?", run.ID, "running", run.LeaseOwner).Update("lease_expires_at", lease).Error
|
||||
func (service *Service) renewAutoMatchRun(ctx context.Context, run models.ShopeeSpecAutoMatchRun) error {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
return service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := lockAutoMatchRun(ctx, tx, run); err != nil {
|
||||
return err
|
||||
}
|
||||
result := autoMatchOwnedRun(tx, run).Update("lease_expires_at", time.Now().UTC().Add(autoMatchLeaseDuration))
|
||||
if result.Error == nil && result.RowsAffected == 0 {
|
||||
// MySQL's changed-row count can be zero when datetime precision
|
||||
// rounds a rapid renewal to the stored value. Under the same row
|
||||
// lock, distinguish that no-op from an expired or lost lease.
|
||||
return lockAutoMatchRun(ctx, tx, run)
|
||||
}
|
||||
return autoMatchOwnedUpdate(result)
|
||||
})
|
||||
}
|
||||
|
||||
func (service *Service) finishAutoMatchRun(run models.ShopeeSpecAutoMatchRun, status string, scanned, eligible, processed, confirmed, unmatched, failed int, summary string) error {
|
||||
func (service *Service) finishAutoMatchRun(ctx context.Context, run models.ShopeeSpecAutoMatchRun, stats autoMatchBatchStats, status, reason string, cursor *uint64) error {
|
||||
now := time.Now().UTC()
|
||||
updates := map[string]any{"status": status, "active_slot": nil, "lease_owner": "", "lease_expires_at": nil, "scanned_count": scanned, "eligible_count": eligible, "processed_count": processed, "confirmed_count": confirmed, "unmatched_count": unmatched, "failed_count": failed, "error_summary": truncateBatchText(summary), "finished_at": now}
|
||||
return service.DB.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ? AND status = ? AND lease_owner = ?", run.ID, "running", run.LeaseOwner).Updates(updates).Error
|
||||
updates := map[string]any{"status": status, "active_slot": nil, "lease_owner": "", "lease_expires_at": nil, "scanned_count": stats.scanned, "eligible_count": stats.eligible, "processed_count": stats.processed, "confirmed_count": stats.confirmed, "unmatched_count": stats.unmatched, "failed_count": stats.failed, "error_summary": truncateBatchText(stats.summary), "finished_at": now, "resume_after_id": cursor, "stop_reason": reason}
|
||||
return service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := lockAutoMatchRun(ctx, tx, run); err != nil {
|
||||
return err
|
||||
}
|
||||
return autoMatchOwnedUpdate(autoMatchOwnedRun(tx, run).Updates(updates))
|
||||
})
|
||||
}
|
||||
|
||||
var errAutoMatchLeaseLost = errors.New("auto match run lease lost")
|
||||
|
||||
type autoMatchRunContextKey struct{}
|
||||
type autoMatchRunGuard struct {
|
||||
run models.ShopeeSpecAutoMatchRun
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
|
||||
func autoMatchOwnedRun(db *gorm.DB, run models.ShopeeSpecAutoMatchRun) *gorm.DB {
|
||||
return db.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ? AND status = ? AND active_slot = ? AND lease_owner = ? AND lease_expires_at > ?", run.ID, "running", 1, run.LeaseOwner, time.Now().UTC())
|
||||
}
|
||||
|
||||
func autoMatchOwnedUpdate(result *gorm.DB) error {
|
||||
if result.Error != nil {
|
||||
return result.Error
|
||||
}
|
||||
if result.RowsAffected != 1 {
|
||||
return errAutoMatchLeaseLost
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func lockAutoMatchRun(ctx context.Context, tx *gorm.DB, run models.ShopeeSpecAutoMatchRun) error {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
var owned models.ShopeeSpecAutoMatchRun
|
||||
err := autoMatchOwnedRun(tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}), run).Take(&owned).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return errAutoMatchLeaseLost
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
if owned.LeaseExpiresAt == nil || !owned.LeaseExpiresAt.After(time.Now().UTC()) {
|
||||
return errAutoMatchLeaseLost
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Only scheduled/manual batch runs attach this context; individual matching
|
||||
// keeps its existing behavior. Reuse the caller's transaction for row locks.
|
||||
func checkAutoMatchRunContext(ctx context.Context, db *gorm.DB) error {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
guard, ok := ctx.Value(autoMatchRunContextKey{}).(autoMatchRunGuard)
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
var owned models.ShopeeSpecAutoMatchRun
|
||||
err := autoMatchOwnedRun(db.WithContext(ctx), guard.run).Take(&owned).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
err = errAutoMatchLeaseLost
|
||||
}
|
||||
if err == nil && (owned.LeaseExpiresAt == nil || !owned.LeaseExpiresAt.After(time.Now().UTC())) {
|
||||
err = errAutoMatchLeaseLost
|
||||
}
|
||||
if errors.Is(err, errAutoMatchLeaseLost) {
|
||||
guard.cancel()
|
||||
}
|
||||
if err == nil {
|
||||
err = ctx.Err()
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func autoMatchWorkSkip(work models.ShopeeSpecAutoMatchWorkItem, fingerprint string, now time.Time) string {
|
||||
if work.InputFingerprint != fingerprint {
|
||||
return ""
|
||||
}
|
||||
if work.Status == "completed" || work.Status == "unmatched" {
|
||||
return "unchanged"
|
||||
}
|
||||
if work.AttemptCount >= maxAutoMatchAttempts {
|
||||
return "max_retry"
|
||||
}
|
||||
if work.NextAttemptAt != nil && work.NextAttemptAt.After(now) {
|
||||
return "cooldown"
|
||||
}
|
||||
if work.Status == "running" && work.LeaseExpiresAt != nil && work.LeaseExpiresAt.After(now) {
|
||||
return "occupied"
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
type autoMatchBatchStats struct {
|
||||
scanned, eligible, processed, confirmed, unmatched, failed int
|
||||
summary string
|
||||
skips map[string]int
|
||||
checkpointLoaded bool
|
||||
}
|
||||
|
||||
func (stats *autoMatchBatchStats) skip(reason string) {
|
||||
if stats.skips == nil {
|
||||
stats.skips = map[string]int{}
|
||||
}
|
||||
stats.skips[reason]++
|
||||
}
|
||||
|
||||
func (service *Service) endAutoMatchBatch(ctx context.Context, run models.ShopeeSpecAutoMatchRun, stats autoMatchBatchStats, cursor uint64, reason string, cause error) error {
|
||||
checkpoint := &cursor
|
||||
if !stats.checkpointLoaded {
|
||||
checkpoint = nil
|
||||
}
|
||||
status := "completed"
|
||||
if stats.failed > 0 {
|
||||
status = "completed_partial"
|
||||
}
|
||||
if errors.Is(cause, errAutoMatchLeaseLost) || errors.Is(context.Cause(ctx), errAutoMatchLeaseLost) {
|
||||
reason, cause = "lease_lost", errAutoMatchLeaseLost
|
||||
} else if ctx.Err() != nil {
|
||||
reason = "time_budget"
|
||||
}
|
||||
if reason == "lease_lost" || reason == "error" {
|
||||
checkpoint = nil
|
||||
status = "failed"
|
||||
stats.failed++
|
||||
stats.summary = safeBatchError(cause)
|
||||
}
|
||||
// Finalization is the sole exception to the scan deadline: a fresh bounded
|
||||
// context records the last fully decided item after a time-budget stop.
|
||||
finishCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
var finishErr error
|
||||
if reason != "lease_lost" {
|
||||
finishErr = service.finishAutoMatchRun(finishCtx, run, stats, status, reason, checkpoint)
|
||||
}
|
||||
if errors.Is(finishErr, errAutoMatchLeaseLost) {
|
||||
reason = "lease_lost"
|
||||
} else if finishErr != nil {
|
||||
reason = "error"
|
||||
}
|
||||
log.Infof("shopee_spec_auto_match run_id=%d stop_reason=%s scanned=%d eligible=%d processed=%d confirmed=%d unmatched=%d failed=%d skip_no_specs=%d skip_confirmed=%d skip_unchanged=%d skip_max_retry=%d skip_cooldown=%d skip_occupied=%d", run.ID, reason, stats.scanned, stats.eligible, stats.processed, stats.confirmed, stats.unmatched, stats.failed, stats.skips["no_specs"], stats.skips["confirmed"], stats.skips["unchanged"], stats.skips["max_retry"], stats.skips["cooldown"], stats.skips["occupied"])
|
||||
if finishErr != nil {
|
||||
return finishErr
|
||||
}
|
||||
if reason == "time_budget" {
|
||||
return nil
|
||||
}
|
||||
return cause
|
||||
}
|
||||
|
||||
func batchErrorCode(err error) string {
|
||||
|
||||
@@ -3,6 +3,7 @@ package shopeeproduct
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"go-admin/app/goauto/models"
|
||||
|
||||
@@ -66,7 +67,8 @@ func TestUnchangedUnmatchedWorkIsNotClaimedAgain(t *testing.T) {
|
||||
db := openTestDB(t)
|
||||
service := NewService(db)
|
||||
one := uint8(1)
|
||||
run := models.ShopeeSpecAutoMatchRun{RequestID: uuid.NewString(), Trigger: "manual", Status: "running", ActiveSlot: &one, LeaseOwner: uuid.NewString(), BatchLimit: 20}
|
||||
lease := time.Now().UTC().Add(autoMatchLeaseDuration)
|
||||
run := models.ShopeeSpecAutoMatchRun{RequestID: uuid.NewString(), Trigger: "manual", Status: "running", ActiveSlot: &one, LeaseOwner: uuid.NewString(), LeaseExpiresAt: &lease, BatchLimit: 20}
|
||||
if err := db.Create(&run).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -74,7 +76,9 @@ func TestUnchangedUnmatchedWorkIsNotClaimedAgain(t *testing.T) {
|
||||
if err != nil || !claimed {
|
||||
t.Fatalf("work=%+v claimed=%v err=%v", work, claimed, err)
|
||||
}
|
||||
service.completeAutoMatchWork(work, "fingerprint", 0, 1, nil)
|
||||
if err := service.completeAutoMatchWork(context.Background(), work, "fingerprint", 0, 1, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, claimed, err = service.claimAutoMatchWork(context.Background(), run, 99, "fingerprint")
|
||||
if err != nil || claimed {
|
||||
t.Fatalf("unchanged unmatched claimed=%v err=%v", claimed, err)
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
package shopeeproduct
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"gorm.io/driver/mysql"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
// Opt-in, synthetic SELECTs only: no schema selection is required, and no
|
||||
// tables, production rows, migrations or credentials are written or logged.
|
||||
func TestAutoMatchScanMySQLJSONCoarseFilter(t *testing.T) {
|
||||
dsn := os.Getenv("GOAUTO_TEST_MYSQL_READONLY_DSN")
|
||||
if dsn == "" {
|
||||
t.Skip("set GOAUTO_TEST_MYSQL_READONLY_DSN to opt in to read-only MySQL compatibility checks")
|
||||
}
|
||||
db, err := gorm.Open(mysql.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
if err != nil {
|
||||
t.Fatal("cannot connect to opted-in MySQL")
|
||||
}
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
t.Fatal("cannot access opted-in MySQL connection")
|
||||
}
|
||||
t.Cleanup(func() { sqlDB.Close() })
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
for _, tc := range []struct {
|
||||
name, json string
|
||||
want int
|
||||
}{
|
||||
{"empty_string", `""`, 0}, {"empty_array", `[]`, 0}, {"json_null", `null`, 0}, {"spaced_array", `[ ]`, 0},
|
||||
{"empty_values", `[{"role":"size","values":[]}]`, 1}, {"size_only", sizeScanSpecs, 1},
|
||||
{"color_only", `[{"role":"color","values":[{"name":"黑色"}]}]`, 1},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
var count int
|
||||
if err := db.WithContext(ctx).Raw("SELECT COUNT(*) FROM (SELECT CAST(? AS JSON) AS specs_json) shopee_product WHERE "+autoMatchNonEmptySpecsSQL, tc.json, "", "[]", "null", `""`).Scan(&count).Error; err != nil {
|
||||
t.Fatal("MySQL JSON coarse filter query failed")
|
||||
}
|
||||
if count != tc.want {
|
||||
t.Fatalf("count=%d want=%d", count, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,784 @@
|
||||
package shopeeproduct
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"go-admin/app/goauto/models"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const sizeScanSpecs = `[{"name":"尺码","role":"size","values":[{"name":"XL","source":"import"}]}]`
|
||||
const otherScanSpecs = `[{"name":"材质","role":"other","values":[{"name":"棉","source":"import"}]}]`
|
||||
|
||||
func openScanTestDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
db := openTestDB(t)
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Closing the final connection releases the named in-memory database,
|
||||
// including when go test repeats the same t.Name via -count.
|
||||
t.Cleanup(func() {
|
||||
if err := sqlDB.Close(); err != nil {
|
||||
t.Errorf("close scan test database: %v", err)
|
||||
}
|
||||
})
|
||||
return db
|
||||
}
|
||||
|
||||
func seedScanProducts(t *testing.T, db *gorm.DB, pddID uint64, count int, specs string) []models.ShopeeProduct {
|
||||
t.Helper()
|
||||
products := make([]models.ShopeeProduct, count)
|
||||
for i := range products {
|
||||
products[i] = models.ShopeeProduct{ShopeeItemID: uuid.NewString(), PDDProductID: &pddID, SpecsJSON: specs}
|
||||
}
|
||||
if err := db.CreateInBatches(&products, 100).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return products
|
||||
}
|
||||
|
||||
func runScan(t *testing.T, service *Service, limit int) *AutoMatchRunView {
|
||||
t.Helper()
|
||||
run, created, err := service.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, limit)
|
||||
if err != nil || !created {
|
||||
t.Fatalf("start: created=%v err=%v", created, err)
|
||||
}
|
||||
if err := service.ProcessAutoMatchRun(context.Background(), run.ID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
latest, err := service.LatestAutoMatchRun(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return latest
|
||||
}
|
||||
|
||||
func scanCheckpoint(t *testing.T, run *AutoMatchRunView, cursor uint64, reason string) {
|
||||
t.Helper()
|
||||
raw, err := json.Marshal(run)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var fields map[string]any
|
||||
if err := json.Unmarshal(raw, &fields); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if fields["resumeAfterId"] != float64(cursor) || fields["stopReason"] != reason {
|
||||
t.Fatalf("checkpoint got cursor=%v reason=%v; want %d %s", fields["resumeAfterId"], fields["stopReason"], cursor, reason)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanPassesLongEmptyPrefix(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
seedScanProducts(t, db, pdd.ID, 1812, `[]`)
|
||||
seedScanProducts(t, db, pdd.ID, 1, sizeScanSpecs)
|
||||
run := runScan(t, NewService(db), 20)
|
||||
if run.ProcessedCount != 1 || run.ScannedCount != 1 || run.ConfirmedCount != 1 {
|
||||
t.Fatalf("run=%+v", run)
|
||||
}
|
||||
scanCheckpoint(t, run, 0, "end_of_scan")
|
||||
}
|
||||
|
||||
func TestAutoMatchScanRotatesAcrossServiceRestart(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
prefix := seedScanProducts(t, db, pdd.ID, 2001, otherScanSpecs)
|
||||
seedScanProducts(t, db, pdd.ID, 1, sizeScanSpecs)
|
||||
first := runScan(t, NewService(db), 20)
|
||||
if first.ScannedCount != 2000 || first.ProcessedCount != 0 {
|
||||
t.Fatalf("first=%+v", first)
|
||||
}
|
||||
scanCheckpoint(t, first, prefix[1999].ID, "scan_budget")
|
||||
second := runScan(t, NewService(db), 20)
|
||||
if second.ScannedCount != 2 || second.ProcessedCount != 1 {
|
||||
t.Fatalf("second=%+v", second)
|
||||
}
|
||||
scanCheckpoint(t, second, 0, "end_of_scan")
|
||||
if err := db.Model(&prefix[0]).Update("specs_json", sizeScanSpecs).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
third := runScan(t, NewService(db), 20)
|
||||
if third.ProcessedCount != 1 {
|
||||
t.Fatalf("changed low ID not visited: %+v", third)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanBatchLimitKeepsLastExaminedOnShortPage(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
products := seedScanProducts(t, db, pdd.ID, 3, sizeScanSpecs)
|
||||
for i := 0; i < 3; i++ {
|
||||
run := runScan(t, NewService(db), 1)
|
||||
if run.ProcessedCount != 1 || run.ScannedCount != 1 {
|
||||
t.Fatalf("run=%+v", run)
|
||||
}
|
||||
if i < 2 {
|
||||
scanCheckpoint(t, run, products[i].ID, "batch_limit")
|
||||
} else {
|
||||
scanCheckpoint(t, run, 0, "end_of_scan")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanExactBudgetDoesNotAssumeEnd(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
products := seedScanProducts(t, db, pdd.ID, 2000, otherScanSpecs)
|
||||
first := runScan(t, NewService(db), 20)
|
||||
scanCheckpoint(t, first, products[1999].ID, "scan_budget")
|
||||
second := runScan(t, NewService(db), 20)
|
||||
scanCheckpoint(t, second, 0, "end_of_scan")
|
||||
if second.ScannedCount != 0 {
|
||||
t.Fatalf("second scanned %d", second.ScannedCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchEligibilityRejectsUnusableSpecsBeforePDDRead(t *testing.T) {
|
||||
for i, specs := range []string{"", `[]`, `null`, `""`, otherScanSpecs, `[{"role":"size","values":[]}]`, `[{"role":"color","values":[{"name":" "}]}]`} {
|
||||
t.Run(fmt.Sprint(i), func(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
missing := uint64(999)
|
||||
_, eligible, err := NewService(db).autoMatchEligibility(context.Background(), models.ShopeeProduct{PDDProductID: &missing, SpecsJSON: specs})
|
||||
if err != nil || eligible {
|
||||
t.Fatalf("eligible=%v err=%v", eligible, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanEmptyRepresentationsAndSingleDimension(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
for _, specs := range []string{`[]`, `null`, `""`} {
|
||||
seedScanProducts(t, db, pdd.ID, 1, specs)
|
||||
}
|
||||
seedScanProducts(t, db, pdd.ID, 1, `[{"name":"颜色","role":"color","values":[{"name":"黑色","source":"import"}]}]`)
|
||||
seedScanProducts(t, db, pdd.ID, 1, sizeScanSpecs)
|
||||
run := runScan(t, NewService(db), 20)
|
||||
if run.ScannedCount != 2 || run.ProcessedCount != 2 || run.ConfirmedCount != 2 {
|
||||
t.Fatalf("run=%+v", run)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanLatestCommittedZeroWinsAndNullIsIgnored(t *testing.T) {
|
||||
for _, latest := range []uint64{0, 2} {
|
||||
t.Run(fmt.Sprint(latest), func(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
products := seedScanProducts(t, db, pdd.ID, 3, otherScanSpecs)
|
||||
old := uint64(1)
|
||||
for _, checkpoint := range []*uint64{&old, &latest, nil} {
|
||||
run := models.ShopeeSpecAutoMatchRun{RequestID: uuid.NewString(), Trigger: "manual", Status: "completed", ResumeAfterID: checkpoint}
|
||||
if err := db.Create(&run).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
run := runScan(t, NewService(db), 20)
|
||||
if run.ScannedCount != len(products)-int(latest) {
|
||||
t.Fatalf("wrong checkpoint: %+v", run)
|
||||
}
|
||||
scanCheckpoint(t, run, 0, "end_of_scan")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchRenewAndFinishRejectLostLease(t *testing.T) {
|
||||
for _, change := range []string{"owner", "expired", "slot", "status"} {
|
||||
t.Run(change, func(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
updates := map[string]any{}
|
||||
switch change {
|
||||
case "owner":
|
||||
updates["lease_owner"] = "new-owner"
|
||||
case "expired":
|
||||
updates["lease_expires_at"] = time.Now().UTC().Add(-time.Second)
|
||||
case "slot":
|
||||
updates["active_slot"] = nil
|
||||
case "status":
|
||||
updates["status"] = "failed"
|
||||
}
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ?", run.ID).Updates(updates).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.renewAutoMatchRun(context.Background(), run.ShopeeSpecAutoMatchRun); err != errAutoMatchLeaseLost {
|
||||
t.Fatalf("renew=%v", err)
|
||||
}
|
||||
cursor := uint64(999)
|
||||
if err := s.finishAutoMatchRun(context.Background(), run.ShopeeSpecAutoMatchRun, autoMatchBatchStats{}, "completed", "end_of_scan", &cursor); err != errAutoMatchLeaseLost {
|
||||
t.Fatalf("finish=%v", err)
|
||||
}
|
||||
var current models.ShopeeSpecAutoMatchRun
|
||||
if err := db.First(¤t, run.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if current.ResumeAfterID != nil || current.StopReason != "" {
|
||||
t.Fatalf("old owner committed: %+v", current)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanDatabaseErrorDoesNotCommitCheckpoint(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
products := seedScanProducts(t, db, pdd.ID, 2, sizeScanSpecs)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Callback().Row().Before("gorm:row").Register("test_settings_error", func(tx *gorm.DB) {
|
||||
if tx.Statement.Table == "ai_matching_setting" {
|
||||
tx.AddError(fmt.Errorf("synthetic database error"))
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.ProcessAutoMatchRun(context.Background(), run.ID); err == nil {
|
||||
t.Fatal("database error ignored")
|
||||
}
|
||||
if err := db.Callback().Row().Remove("test_settings_error"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
latest, err := s.LatestAutoMatchRun(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if latest.Status != "failed" || latest.StopReason != "error" || latest.ResumeAfterID != nil {
|
||||
t.Fatalf("latest=%+v", latest)
|
||||
}
|
||||
next := runScan(t, NewService(db), 20)
|
||||
if next.ProcessedCount != len(products) {
|
||||
t.Fatalf("restart=%+v", next)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchClaimRejectsLostOrExpiredRun(t *testing.T) {
|
||||
for _, change := range []string{"owner", "expired", "slot", "status"} {
|
||||
t.Run(change, func(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
updates := map[string]any{}
|
||||
switch change {
|
||||
case "owner":
|
||||
updates["lease_owner"] = uuid.NewString()
|
||||
case "expired":
|
||||
updates["lease_expires_at"] = time.Now().UTC().Add(-time.Second)
|
||||
case "slot":
|
||||
updates["active_slot"] = nil
|
||||
case "status":
|
||||
updates["status"] = "failed"
|
||||
}
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ?", run.ID).Updates(updates).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, claimed, err := s.claimAutoMatchWork(context.Background(), run.ShopeeSpecAutoMatchRun, 999, "fingerprint")
|
||||
if err == nil || claimed {
|
||||
t.Fatalf("lost run claimed=%v err=%v", claimed, err)
|
||||
}
|
||||
var count int64
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchWorkItem{}).Count(&count).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if count != 0 {
|
||||
t.Fatalf("lost run wrote work: %d", count)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanOwnerLossDuringProviderStopsNextCallAndSave(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
products := seedScanProducts(t, db, pdd.ID, 2, `[{"name":"颜色","role":"color","values":[{"name":"深黑","source":"import"}]},{"name":"尺码","role":"size","values":[{"name":"大号","source":"import"}]}]`)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var calls atomic.Int32
|
||||
provider := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
calls.Add(1)
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ?", run.ID).Update("lease_owner", "replacement-owner").Error; err != nil {
|
||||
t.Error(err)
|
||||
}
|
||||
chatCompletionResponder(`{"suggestions":[{"sourceId":"s1","candidateId":"c1","confidence":0.96,"reason":"unique match"}]}`)(w, r)
|
||||
}))
|
||||
defer provider.Close()
|
||||
setting := seedEnabledAISetting(t, provider.URL, 0.9)
|
||||
if err := db.Create(&setting).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.ProcessAutoMatchRun(context.Background(), run.ID); err == nil {
|
||||
t.Fatal("owner loss must be returned")
|
||||
}
|
||||
if calls.Load() != 1 {
|
||||
t.Fatalf("provider calls after owner loss: %d", calls.Load())
|
||||
}
|
||||
var current models.ShopeeSpecAutoMatchRun
|
||||
if err := db.First(¤t, run.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if current.Status != "running" || current.LeaseOwner != "replacement-owner" {
|
||||
t.Fatalf("old owner overwrote run: %+v", current)
|
||||
}
|
||||
var product models.ShopeeProduct
|
||||
if err := db.First(&product, products[0].ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if product.SpecsJSON != products[0].SpecsJSON {
|
||||
t.Fatal("old owner saved mapping")
|
||||
}
|
||||
var count int64
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchWorkItem{}).Count(&count).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Fatalf("unexpected work claims: %d", count)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanDeadlineStopsProviderAndKeepsLastDecision(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
first := seedScanProducts(t, db, pdd.ID, 1, otherScanSpecs)[0]
|
||||
products := seedScanProducts(t, db, pdd.ID, 2, `[{"name":"颜色","role":"color","values":[{"name":"深黑","source":"import"}]}]`)
|
||||
var calls atomic.Int32
|
||||
provider := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
calls.Add(1)
|
||||
select {
|
||||
case <-r.Context().Done():
|
||||
case <-time.After(time.Second):
|
||||
}
|
||||
}))
|
||||
defer provider.Close()
|
||||
setting := seedEnabledAISetting(t, provider.URL, 0.9)
|
||||
if err := db.Create(&setting).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
|
||||
defer cancel()
|
||||
started := time.Now()
|
||||
if err := s.ProcessAutoMatchRun(ctx, run.ID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if time.Since(started) > time.Second {
|
||||
t.Fatal("provider outlived batch deadline")
|
||||
}
|
||||
if calls.Load() != 1 {
|
||||
t.Fatalf("provider calls=%d", calls.Load())
|
||||
}
|
||||
latest, err := s.LatestAutoMatchRun(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
scanCheckpoint(t, latest, first.ID, "time_budget")
|
||||
if latest.ScannedCount != 2 || latest.ProcessedCount != 1 || latest.FailedCount != 0 {
|
||||
t.Fatalf("latest=%+v", latest)
|
||||
}
|
||||
var current models.ShopeeProduct
|
||||
if err := db.First(¤t, products[0].ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if current.SpecsJSON != products[0].SpecsJSON {
|
||||
t.Fatal("timeout saved mapping")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanDeadlineBoundsDatabaseAndNoClaimAfterBudget(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
seedScanProducts(t, db, pdd.ID, 2, sizeScanSpecs)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
||||
defer cancel()
|
||||
queries := 0
|
||||
if err := db.Callback().Query().Before("gorm:query").Register("test_database_deadline", func(tx *gorm.DB) {
|
||||
deadline, ok := tx.Statement.Context.Deadline()
|
||||
if !ok || time.Until(deadline) > autoMatchTimeBudget {
|
||||
t.Error("database missed total deadline")
|
||||
}
|
||||
if tx.Statement.Table == "pdd_product" {
|
||||
queries++
|
||||
<-tx.Statement.Context.Done()
|
||||
tx.AddError(tx.Statement.Context.Err())
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.ProcessAutoMatchRun(ctx, run.ID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Callback().Query().Remove("test_database_deadline"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
latest, err := s.LatestAutoMatchRun(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
scanCheckpoint(t, latest, 0, "time_budget")
|
||||
if queries != 1 || latest.ProcessedCount != 0 || latest.ScannedCount != 1 {
|
||||
t.Fatalf("queries=%d latest=%+v", queries, latest)
|
||||
}
|
||||
var count int64
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchWorkItem{}).Count(&count).Error; err != nil || count != 0 {
|
||||
t.Fatalf("work=%d err=%v", count, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanRecoveryIgnoresExpiredCheckpoint(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ?", run.ID).Update("lease_expires_at", time.Now().UTC().Add(-time.Second)).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, created, err := s.StartAutoMatchRun(context.Background(), "scheduled", uuid.NewString(), nil, 20); err != nil || !created {
|
||||
t.Fatalf("recovery created=%v err=%v", created, err)
|
||||
}
|
||||
var old models.ShopeeSpecAutoMatchRun
|
||||
if err := db.First(&old, run.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if old.Status != "failed" || old.StopReason != "lease_lost" || old.ResumeAfterID != nil {
|
||||
t.Fatalf("old=%+v", old)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchScanExactPageBoundaryAndCandidateFilters(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
seedScanProducts(t, db, pdd.ID, 200, otherScanSpecs)
|
||||
seedScanProducts(t, db, pdd.ID, 1, sizeScanSpecs)
|
||||
disabled := seedPDDProduct(t, db, "disabled")
|
||||
seedScanProducts(t, db, disabled.ID, 1, sizeScanSpecs)
|
||||
deleted := seedScanProducts(t, db, pdd.ID, 1, sizeScanSpecs)[0]
|
||||
if err := db.Delete(&deleted).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
unlinked := models.ShopeeProduct{ShopeeItemID: uuid.NewString(), SpecsJSON: sizeScanSpecs}
|
||||
if err := db.Create(&unlinked).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
run := runScan(t, NewService(db), 20)
|
||||
if run.ScannedCount != 201 || run.ProcessedCount != 1 {
|
||||
t.Fatalf("run=%+v", run)
|
||||
}
|
||||
scanCheckpoint(t, run, 0, "end_of_scan")
|
||||
}
|
||||
|
||||
func TestAutoMatchWorkRetryAndCooldownPreserved(t *testing.T) {
|
||||
for _, status := range []string{"completed", "unmatched", "max_retry", "cooldown", "occupied", "retryable"} {
|
||||
t.Run(status, func(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
future := time.Now().UTC().Add(time.Hour)
|
||||
work := models.ShopeeSpecAutoMatchWorkItem{ShopeeProductID: 99, InputFingerprint: "same", Status: status, AttemptCount: 1}
|
||||
switch status {
|
||||
case "max_retry":
|
||||
work.Status, work.AttemptCount = "failed", 3
|
||||
case "cooldown":
|
||||
work.Status, work.NextAttemptAt = "failed", &future
|
||||
case "occupied":
|
||||
work.Status, work.LeaseExpiresAt = "running", &future
|
||||
case "retryable":
|
||||
work.Status = "failed"
|
||||
}
|
||||
if err := db.Create(&work).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, claimed, err := s.claimAutoMatchWork(context.Background(), run.ShopeeSpecAutoMatchRun, 99, "same")
|
||||
if err != nil || claimed != (status == "retryable") {
|
||||
t.Fatalf("claimed=%v err=%v", claimed, err)
|
||||
}
|
||||
if claimed {
|
||||
if got.AttemptCount != 2 {
|
||||
t.Fatalf("attempts=%d", got.AttemptCount)
|
||||
}
|
||||
if err := s.completeAutoMatchWork(context.Background(), got, "same", 0, 0, aiUnavailable("synthetic unavailable")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var saved models.ShopeeSpecAutoMatchWorkItem
|
||||
if err := db.First(&saved, got.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if saved.Status != "failed" || saved.NextAttemptAt == nil || saved.LastErrorCode != CodeAIUnavailable {
|
||||
t.Fatalf("saved=%+v", saved)
|
||||
}
|
||||
}
|
||||
changed, claimed, err := s.claimAutoMatchWork(context.Background(), run.ShopeeSpecAutoMatchRun, 99, "changed")
|
||||
if err != nil || !claimed || changed.AttemptCount != 1 {
|
||||
t.Fatalf("changed=%+v claimed=%v err=%v", changed, claimed, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchClaimDoesNotUseLeaseTimeBeforeLockWait(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
expiry := time.Now().UTC().Add(50 * time.Millisecond)
|
||||
if err := db.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ?", run.ID).Update("lease_expires_at", expiry).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Callback().Query().Before("gorm:query").Register("test_lock_wait", func(tx *gorm.DB) {
|
||||
if tx.Statement.Table == "shopee_spec_auto_match_run" {
|
||||
time.Sleep(time.Until(expiry) + 10*time.Millisecond)
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, claimed, err := s.claimAutoMatchWork(context.Background(), run.ShopeeSpecAutoMatchRun, 99, "same")
|
||||
if err != errAutoMatchLeaseLost || claimed {
|
||||
t.Fatalf("claimed=%v err=%v", claimed, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchProviderGuardDatabaseErrorMustNotRetryOrCheckpoint(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
seedScanProducts(t, db, pdd.ID, 1, `[{"name":"颜色","role":"color","values":[{"name":"深黑","source":"import"}]}]`)
|
||||
var calls atomic.Int32
|
||||
provider := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
calls.Add(1)
|
||||
chatCompletionResponder(`{"suggestions":[{"sourceId":"s1","candidateId":"c1","confidence":0.96,"reason":"unique match"}]}`)(w, r)
|
||||
}))
|
||||
defer provider.Close()
|
||||
setting := seedEnabledAISetting(t, provider.URL, 0.9)
|
||||
if err := db.Create(&setting).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// The settings read in suggestMappings immediately precedes its provider guard.
|
||||
armed, injected := false, false
|
||||
if err := db.Callback().Query().Before("gorm:query").Register("test_provider_guard_error", func(tx *gorm.DB) {
|
||||
if tx.Statement.Table == "ai_matching_setting" {
|
||||
armed = true
|
||||
}
|
||||
if armed && !injected && tx.Statement.Table == "shopee_spec_auto_match_run" {
|
||||
injected = true
|
||||
tx.AddError(fmt.Errorf("synthetic provider guard database error"))
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = s.ProcessAutoMatchRun(context.Background(), run.ID)
|
||||
if err == nil || !injected || calls.Load() != 0 {
|
||||
t.Fatalf("err=%v injected=%v calls=%d", err, injected, calls.Load())
|
||||
}
|
||||
latest, err := s.LatestAutoMatchRun(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if latest.ResumeAfterID != nil || latest.StopReason != "error" {
|
||||
t.Fatalf("latest=%+v", latest)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchProviderSettingsDatabaseErrorMustNotRetryOrCheckpoint(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
pdd := seedPDDProduct(t, db, "active")
|
||||
seedScanProducts(t, db, pdd.ID, 1, `[{"name":"颜色","role":"color","values":[{"name":"深黑","source":"import"}]}]`)
|
||||
var calls atomic.Int32
|
||||
provider := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
calls.Add(1)
|
||||
chatCompletionResponder(`{"suggestions":[{"sourceId":"s1","candidateId":"c1","confidence":0.96,"reason":"unique match"}]}`)(w, r)
|
||||
}))
|
||||
defer provider.Close()
|
||||
setting := seedEnabledAISetting(t, provider.URL, 0.9)
|
||||
if err := db.Create(&setting).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
settingReads := 0
|
||||
if err := db.Callback().Query().Before("gorm:query").Register("test_provider_settings_error", func(tx *gorm.DB) {
|
||||
if tx.Statement.Table == "ai_matching_setting" {
|
||||
settingReads++
|
||||
if settingReads == 2 {
|
||||
tx.AddError(fmt.Errorf("synthetic nested settings error"))
|
||||
}
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = s.ProcessAutoMatchRun(context.Background(), run.ID)
|
||||
if err == nil || settingReads != 2 || calls.Load() != 0 {
|
||||
t.Fatalf("err=%v settings_reads=%d calls=%d", err, settingReads, calls.Load())
|
||||
}
|
||||
latest, err := s.LatestAutoMatchRun(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if latest.ResumeAfterID != nil || latest.StopReason != "error" {
|
||||
t.Fatalf("latest=%+v", latest)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchRunDatabaseWriteFailuresAreReturned(t *testing.T) {
|
||||
for _, operation := range []string{"renew", "finish", "work"} {
|
||||
t.Run(operation, func(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
work, _, err := s.claimAutoMatchWork(context.Background(), run.ShopeeSpecAutoMatchRun, 99, "same")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
injected := fmt.Errorf("synthetic update failure")
|
||||
if err := db.Callback().Update().Before("gorm:update").Register("test_update_error", func(tx *gorm.DB) { tx.AddError(injected) }); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
switch operation {
|
||||
case "renew":
|
||||
err = s.renewAutoMatchRun(context.Background(), run.ShopeeSpecAutoMatchRun)
|
||||
case "finish":
|
||||
cursor := uint64(99)
|
||||
err = s.finishAutoMatchRun(context.Background(), run.ShopeeSpecAutoMatchRun, autoMatchBatchStats{}, "completed", "end_of_scan", &cursor)
|
||||
case "work":
|
||||
err = s.completeAutoMatchWork(context.Background(), work, "same", 1, 0, nil)
|
||||
}
|
||||
if err != injected {
|
||||
t.Fatalf("err=%v", err)
|
||||
}
|
||||
var saved models.ShopeeSpecAutoMatchRun
|
||||
if err := db.First(&saved, run.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if saved.ResumeAfterID != nil || saved.Status != "running" {
|
||||
t.Fatalf("saved=%+v", saved)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchDeadlineDuringCheckpointReadCannotCommitFalseHead(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
checkpoint := uint64(1700)
|
||||
old := models.ShopeeSpecAutoMatchRun{RequestID: uuid.NewString(), Trigger: "manual", Status: "completed", ResumeAfterID: &checkpoint}
|
||||
if err := db.Create(&old).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
queries := 0
|
||||
if err := db.Callback().Query().Before("gorm:query").Register("test_checkpoint_timeout", func(tx *gorm.DB) {
|
||||
if tx.Statement.Table == "shopee_spec_auto_match_run" {
|
||||
queries++
|
||||
if queries == 2 {
|
||||
<-tx.Statement.Context.Done()
|
||||
tx.AddError(tx.Statement.Context.Err())
|
||||
}
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
|
||||
defer cancel()
|
||||
if err := s.ProcessAutoMatchRun(ctx, run.ID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
latest, err := s.LatestAutoMatchRun(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if latest.ResumeAfterID != nil || latest.StopReason != "time_budget" {
|
||||
t.Fatalf("unknown cursor committed: %+v", latest)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAutoMatchRenewZeroChangedRowsRequiresLiveOwnership(t *testing.T) {
|
||||
for _, change := range []string{"unchanged", "expired", "owner", "multiple_rows"} {
|
||||
t.Run(change, func(t *testing.T) {
|
||||
db := openScanTestDB(t)
|
||||
s := NewService(db)
|
||||
run, _, err := s.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Callback().Update().After("gorm:update").Register("test_renew_zero_changed", func(tx *gorm.DB) {
|
||||
if tx.Statement.Table != "shopee_spec_auto_match_run" || tx.Error != nil {
|
||||
return
|
||||
}
|
||||
// MySQL reports changed rows by default: datetime(3) may round a
|
||||
// same-millisecond renewal to the value already stored.
|
||||
switch change {
|
||||
case "expired":
|
||||
err = tx.Session(&gorm.Session{NewDB: true}).Exec("UPDATE shopee_spec_auto_match_run SET lease_expires_at = ? WHERE id = ?", time.Now().UTC().Add(-time.Second), run.ID).Error
|
||||
case "owner":
|
||||
err = tx.Session(&gorm.Session{NewDB: true}).Exec("UPDATE shopee_spec_auto_match_run SET lease_owner = ? WHERE id = ?", "replacement-owner", run.ID).Error
|
||||
}
|
||||
if err != nil {
|
||||
tx.AddError(err)
|
||||
}
|
||||
tx.RowsAffected = 0
|
||||
if change == "multiple_rows" {
|
||||
tx.RowsAffected = 2
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = s.renewAutoMatchRun(context.Background(), run.ShopeeSpecAutoMatchRun)
|
||||
if change == "unchanged" {
|
||||
if err != nil {
|
||||
t.Fatalf("live no-op renewal rejected: %v", err)
|
||||
}
|
||||
} else if err != errAutoMatchLeaseLost {
|
||||
t.Fatalf("lost lease accepted after zero changed rows: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
+32
@@ -0,0 +1,32 @@
|
||||
package version_local
|
||||
|
||||
import (
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/cmd/migrate/migration"
|
||||
common "go-admin/common/models"
|
||||
"gorm.io/gorm"
|
||||
"runtime"
|
||||
)
|
||||
|
||||
func init() {
|
||||
_, file, _, _ := runtime.Caller(0)
|
||||
migration.Migrate.SetVersion(migration.GetFilename(file), migrateShopeeSpecAutoMatchResume)
|
||||
}
|
||||
|
||||
func migrateShopeeSpecAutoMatchResume(db *gorm.DB, version string) error {
|
||||
return db.Transaction(func(tx *gorm.DB) error {
|
||||
if !tx.Migrator().HasColumn(&models.ShopeeSpecAutoMatchRun{}, "ResumeAfterID") {
|
||||
// GORM AddColumn omits CHECK tags. Inline the portable constraint so
|
||||
// this stays additive (SQLite otherwise rebuilds tables for checks).
|
||||
if err := tx.Exec("ALTER TABLE shopee_spec_auto_match_run ADD COLUMN resume_after_id BIGINT NULL CONSTRAINT ck_shopee_spec_auto_match_resume CHECK (resume_after_id IS NULL OR resume_after_id >= 0)").Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if !tx.Migrator().HasColumn(&models.ShopeeSpecAutoMatchRun{}, "StopReason") {
|
||||
if err := tx.Migrator().AddColumn(&models.ShopeeSpecAutoMatchRun{}, "StopReason"); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return tx.Where("version = ?", version).FirstOrCreate(&common.Migration{Version: version}).Error
|
||||
})
|
||||
}
|
||||
+79
@@ -0,0 +1,79 @@
|
||||
package version_local
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"go-admin/app/goauto/models"
|
||||
common "go-admin/common/models"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func TestMigrateShopeeSpecAutoMatchResumePreservesLegacyAndIsIdempotent(t *testing.T) {
|
||||
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB, _ := db.DB()
|
||||
t.Cleanup(func() { sqlDB.Close() })
|
||||
if err := db.Exec("CREATE TABLE shopee_spec_auto_match_run (id integer primary key, status varchar(24) NOT NULL)").Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.Exec("INSERT INTO shopee_spec_auto_match_run(id,status) VALUES(1,'completed')").Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.AutoMigrate(&common.Migration{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i := 0; i < 2; i++ {
|
||||
if err := migrateShopeeSpecAutoMatchResume(db, "test_auto_match_resume"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
var row models.ShopeeSpecAutoMatchRun
|
||||
if err := db.First(&row, 1).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if row.Status != "completed" || row.ResumeAfterID != nil || row.StopReason != "" {
|
||||
t.Fatalf("legacy changed: %+v", row)
|
||||
}
|
||||
columns, err := db.Migrator().ColumnTypes(&row)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(columns) != 4 {
|
||||
t.Fatalf("migration must append exactly two columns: %d", len(columns))
|
||||
}
|
||||
if err := db.Model(&row).UpdateColumns(map[string]any{"resume_after_id": 0, "stop_reason": "end_of_scan"}).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := migrateShopeeSpecAutoMatchResume(db, "test_auto_match_resume"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
row = models.ShopeeSpecAutoMatchRun{}
|
||||
if err := db.First(&row, 1).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if row.ResumeAfterID == nil || *row.ResumeAfterID != 0 || row.StopReason != "end_of_scan" {
|
||||
t.Fatalf("zero overwritten: %+v", row)
|
||||
}
|
||||
var count int64
|
||||
if err := db.Model(&common.Migration{}).Count(&count).Error; err != nil || count != 1 {
|
||||
t.Fatalf("versions=%d err=%v", count, err)
|
||||
}
|
||||
raw, err := json.Marshal(row)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var oldClient struct {
|
||||
ID uint64 `json:"id"`
|
||||
Status string `json:"status"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &oldClient); err != nil || oldClient.ID != 1 || oldClient.Status != "completed" {
|
||||
t.Fatalf("old client=%+v err=%v", oldClient, err)
|
||||
}
|
||||
if err := db.Exec("UPDATE shopee_spec_auto_match_run SET resume_after_id = -1 WHERE id = 1").Error; err == nil {
|
||||
t.Fatal("negative checkpoint accepted")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user