Compare commits

...
14 changed files with 1377 additions and 109 deletions
+14 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Architecture-and-Code-Map
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.-
wiki_revision: 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重试策略保持不变。
+13 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Business-Rules-and-Glossary
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.-
wiki_revision: 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是规格项数。商品变化后可能需要等扫描回绕,不保证固定小时内全部处理。
+12 -2
View File
@@ -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,预算退出有持久游标;错误/失租不提交新游标。
+13 -2
View File
@@ -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正式迁移及真实批次验证仍须授权。
+4 -1
View File
@@ -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" }
+29 -3
View File
@@ -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(&current, 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(&current, 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(&current, 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(&current, 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)
}
})
}
}
@@ -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
})
}
@@ -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")
}
}