package repository import ( "database/sql" "errors" "fmt" "strings" "testing" "github.com/go-sql-driver/mysql" _ "modernc.org/sqlite" "cmautobuy/admin/model" ) type skuSpecConflictExecer struct { *sql.DB conflictKey string inserted bool } func (q *skuSpecConflictExecer) Exec(query string, args ...any) (sql.Result, error) { if !q.inserted && strings.HasPrefix(strings.TrimSpace(query), "INSERT INTO shopee_skus") { q.inserted = true _, err := q.DB.Exec(`INSERT INTO shopee_skus( sku_id,goods_id,spec_raw,spec_key,color,size,advice,parse_ok,is_manual, source,field_sources,field_observed_at,source_observed_at,created_at,updated_at ) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, "catalog:existing", "GOODS-1", "白色,M", "白色,M", "目录白", "M", "", 1, 0, "xlsx-processed", `{"color":"xlsx-processed","size":"xlsx-processed","advice":"","sku_code":""}`, `{"color":"2026-08-21T00:00:00Z","size":"2026-08-21T00:00:00Z","advice":"","sku_code":""}`, "2026-08-21T00:00:00Z", "2026-08-21T00:00:00Z", "2026-08-21T00:00:00Z") if err != nil { return nil, err } return nil, &mysql.MySQLError{Number: 1062, Message: "Duplicate entry for key '" + q.conflictKey + "'"} } return q.DB.Exec(query, args...) } func (q *skuSpecConflictExecer) QueryRow(query string, args ...any) *sql.Row { return q.DB.QueryRow(strings.TrimSuffix(query, " FOR UPDATE"), args...) } func newSKUConflictTestDB(t *testing.T, conflictKey string) *skuSpecConflictExecer { t.Helper() db, err := sql.Open("sqlite", ":memory:") if err != nil { t.Fatal(err) } t.Cleanup(func() { db.Close() }) if _, err := db.Exec(`CREATE TABLE shopee_skus ( sku_id TEXT PRIMARY KEY, shopee_sku_id TEXT, goods_id TEXT NOT NULL, spec_raw TEXT NOT NULL, spec_key TEXT NOT NULL, color TEXT, size TEXT, advice TEXT, parse_ok INTEGER NOT NULL, sku_code TEXT, is_manual INTEGER NOT NULL, source TEXT, field_sources TEXT, field_observed_at TEXT, source_observed_at TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, UNIQUE(goods_id, spec_key) )`); err != nil { t.Fatal(err) } return &skuSpecConflictExecer{DB: db, conflictKey: conflictKey} } func TestUpsertCatalogShopeeSKU_规格唯一键并发冲突后安全重读(t *testing.T) { q := newSKUConflictTestDB(t, "uq_shopee_skus_spec") outcome, err := UpsertCatalogShopeeSKU(q, CatalogShopeeSKUInput{ RecordID: "syb:incoming", GoodsID: "GOODS-1", SpecRaw: "白色,M", SpecKey: "白色,M", Color: "顺运宝白", Size: "M", ParseOK: true, Source: "syb", ObservedAt: "2026-08-21T00:01:00Z", Now: "2026-08-21T00:01:00Z", UpdatePolicy: "fill_missing", }) if err != nil || outcome != CatalogSKUSkipped { t.Fatalf("并发撞键后应重读并按既有规则跳过:outcome=%q err=%v", outcome, err) } var id, color, source string if err := q.QueryRow(`SELECT sku_id,color,source FROM shopee_skus WHERE goods_id='GOODS-1' AND spec_key='白色,M'`).Scan(&id, &color, &source); err != nil { t.Fatal(err) } if id != "catalog:existing" || color != "目录白" || source != "xlsx-processed" { t.Fatalf("竞争方目录数据不得被 SYB 覆盖:id=%q color=%q source=%q", id, color, source) } } func TestUpsertCatalogShopeeSKU_非规格唯一键冲突仍失败(t *testing.T) { q := newSKUConflictTestDB(t, "uq_shopee_skus_external") _, err := UpsertCatalogShopeeSKU(q, CatalogShopeeSKUInput{ RecordID: "syb:incoming", GoodsID: "GOODS-1", SpecRaw: "白色,M", SpecKey: "白色,M", Color: "顺运宝白", Size: "M", ParseOK: true, Source: "syb", ObservedAt: "2026-08-21T00:01:00Z", Now: "2026-08-21T00:01:00Z", UpdatePolicy: "fill_missing", }) var mysqlErr *mysql.MySQLError if !errors.As(err, &mysqlErr) || mysqlErr.Number != 1062 { t.Fatalf("非规格唯一键冲突必须原样失败:%v", err) } } func TestCatalogImportRun_登记查询重放和冲突(t *testing.T) { db, err := sql.Open("sqlite", ":memory:") if err != nil { t.Fatal(err) } defer db.Close() if _, err := db.Exec(`CREATE TABLE catalog_import_runs ( source TEXT NOT NULL,batch_id TEXT NOT NULL,request_hash TEXT NOT NULL,status TEXT NOT NULL, request_count INTEGER NOT NULL,conflict_count INTEGER NOT NULL,observed_at TEXT,update_policy TEXT NOT NULL DEFAULT 'fill_missing',last_request_at TEXT NOT NULL, last_conflict_at TEXT,shopee_created INTEGER NOT NULL DEFAULT 0,shopee_updated INTEGER NOT NULL DEFAULT 0,shopee_fields_filled INTEGER NOT NULL DEFAULT 0,shopee_fields_same_source_updated INTEGER NOT NULL DEFAULT 0,shopee_fields_manual_skipped INTEGER NOT NULL DEFAULT 0,shopee_fields_stale_skipped INTEGER NOT NULL DEFAULT 0, sku_created INTEGER NOT NULL DEFAULT 0,sku_updated INTEGER NOT NULL DEFAULT 0,sku_filled INTEGER NOT NULL DEFAULT 0,sku_same_source_updated INTEGER NOT NULL DEFAULT 0,sku_skipped INTEGER NOT NULL DEFAULT 0,sku_manual_skipped INTEGER NOT NULL DEFAULT 0,sku_stale_skipped INTEGER NOT NULL DEFAULT 0,pdd_created INTEGER NOT NULL DEFAULT 0, pdd_updated INTEGER NOT NULL DEFAULT 0,association_created INTEGER NOT NULL DEFAULT 0, association_unchanged INTEGER NOT NULL DEFAULT 0,failure_count INTEGER NOT NULL DEFAULT 0, error_summary TEXT,response_body TEXT,created_at TEXT NOT NULL,finished_at TEXT, PRIMARY KEY(source,batch_id))`); err != nil { t.Fatal(err) } run := model.CatalogImportRun{Source: "script-a", BatchID: "batch-1", RequestHash: "abc", Status: model.CatalogImportProcessing, ObservedAt: "2026-08-11T00:00:00Z", LastRequestAt: "2026-08-11T00:01:00Z", CreatedAt: "2026-08-11T00:01:00Z"} if err := InsertCatalogImportRun(db, run); err != nil { t.Fatal(err) } if err := InsertCatalogImportRun(db, run); !errors.Is(err, ErrCatalogImportRunExists) { t.Fatalf("重复批次应由唯一键拦截:%v", err) } if err := RecordCatalogImportReplay(db, run.Source, run.BatchID, "2026-08-11T00:02:00Z"); err != nil { t.Fatal(err) } if err := RecordCatalogImportConflict(db, run.Source, run.BatchID, "2026-08-11T00:03:00Z"); err != nil { t.Fatal(err) } got, err := GetCatalogImportRun(db, run.Source, run.BatchID) if err != nil { t.Fatal(err) } if got.RequestCount != 3 || got.ConflictCount != 1 || got.LastConflictAt == "" { t.Fatalf("批次计数不正确:%+v", got) } other := run other.Source, other.BatchID = "script-b", "batch-2" if err := InsertCatalogImportRun(db, other); err != nil { t.Fatal(err) } list, total, err := ListCatalogImportRuns(db, "script-a", model.CatalogImportProcessing, 20, 0) if err != nil || total != 1 || len(list) != 1 || list[0].BatchID != "batch-1" { t.Fatalf("分页筛选错误:total=%d list=%+v err=%v", total, list, err) } sources, err := ListCatalogImportSources(db) if err != nil || len(sources) != 2 { t.Fatalf("来源列表错误:%v %v", sources, err) } for i := 3; i <= 1001; i++ { bulk := run bulk.Source = "script-b" bulk.BatchID = fmt.Sprintf("batch-%04d", i) if err := InsertCatalogImportRun(db, bulk); err != nil { t.Fatal(err) } } page, total, err := ListCatalogImportRuns(db, "script-b", "", 20, 20) if err != nil || total != 1000 || len(page) != 20 { t.Fatalf("1000 条数据库分页错误:total=%d rows=%d err=%v", total, len(page), err) } }