diff --git a/server/app/goauto/sybproductfilter/recompute.go b/server/app/goauto/sybproductfilter/recompute.go index 127b6cc..a4af910 100644 --- a/server/app/goauto/sybproductfilter/recompute.go +++ b/server/app/goauto/sybproductfilter/recompute.go @@ -5,7 +5,6 @@ import ( "crypto/sha256" "encoding/hex" "encoding/json" - "fmt" "sort" "strings" "time" @@ -147,25 +146,66 @@ func recomputeVariationSku(rawJSON string) string { return value } -// recomputeFingerprint binds a preview to the exact plan it showed: it is a -// sha256 over the sorted list of "id:direction:ruleId" for every planned -// change, so a genuinely identical plan always produces the same fingerprint -// regardless of query row order (#340 phase 3 review item 2). +// recomputeFingerprintEntry is one change's canonical, unambiguous +// representation for hashing (#340 phase 4 review item 1). It is JSON, not +// naive string concatenation: a naive "id:direction:ruleId" (or any other +// delimiter-joined string) can collide between two different plans whenever +// a field's own text can contain the delimiter or vary in length — e.g. a +// rule keyword containing ":" or newlines could make two distinct plans hash +// identically. encoding/json's field ordering for a fixed struct is stable, +// so this is both deterministic and injective for our purposes. +type recomputeFingerprintEntry struct { + ID uint64 `json:"id"` + Direction string `json:"direction"` + RuleID uint64 `json:"ruleId"` + RuleKind string `json:"ruleKind"` + RuleKeyword string `json:"ruleKeyword"` +} + +// recomputeFingerprint binds a preview to the exact plan it showed, INCLUDING +// the rule evidence that will be written to excluded_rule_kind/ +// excluded_rule_keyword (#340 phase 4 review item 1): two plans that flip the +// exact same id+direction but via a different (or since-edited) rule must +// hash differently, because RecomputeExecute is about to persist exactly +// this rule kind/keyword as this row's excluded_rule_* snapshot — a +// fingerprint that ignored them could let a stale plan through unnoticed +// whenever a rule's keyword/kind changed between preview and execute but the +// set of affected ids/directions happened to stay the same. For the +// excluded_to_pdd direction there is no rule (the row is losing its mark), +// so RuleID/RuleKind/RuleKeyword are left at their zero values, matching what +// gets written (nil/""/""). +// +// It is a sha256 over the JSON-encoded, sorted (by id, then direction) list +// of recomputeFingerprintEntry — sorting the decoded entries themselves +// (not pre-serialized strings) keeps the ordering rule obviously correct +// regardless of how any field is later escaped. func recomputeFingerprint(changes []recomputeChange) string { - entries := make([]string, 0, len(changes)) + entries := make([]recomputeFingerprintEntry, 0, len(changes)) for _, change := range changes { - var ruleID uint64 - if change.ruleID != nil { - ruleID = *change.ruleID - } - direction := "0" + entry := recomputeFingerprintEntry{ID: change.id} if change.toExcluded { - direction = "1" + entry.Direction = DirectionPDDToExcluded + if change.ruleID != nil { + entry.RuleID = *change.ruleID + } + entry.RuleKind = change.ruleKind + entry.RuleKeyword = change.ruleKeyword + } else { + entry.Direction = DirectionExcludedToPDD } - entries = append(entries, fmt.Sprintf("%d:%s:%d", change.id, direction, ruleID)) + entries = append(entries, entry) } - sort.Strings(entries) - sum := sha256.Sum256([]byte(strings.Join(entries, "\n"))) + sort.Slice(entries, func(i, j int) bool { + if entries[i].ID != entries[j].ID { + return entries[i].ID < entries[j].ID + } + return entries[i].Direction < entries[j].Direction + }) + // Marshal errors are impossible here (every field is a plain string/uint64 + // with no cycles), so it is safe to ignore the error and hash whatever + // was produced rather than plumb an error return through every caller. + payload, _ := json.Marshal(entries) + sum := sha256.Sum256(payload) return hex.EncodeToString(sum[:]) } diff --git a/server/app/goauto/sybproductfilter/recompute_mysql_integration_test.go b/server/app/goauto/sybproductfilter/recompute_mysql_integration_test.go new file mode 100644 index 0000000..c4745cd --- /dev/null +++ b/server/app/goauto/sybproductfilter/recompute_mysql_integration_test.go @@ -0,0 +1,351 @@ +package sybproductfilter + +import ( + "context" + "database/sql" + "fmt" + "os" + "regexp" + "testing" + "time" + + "go-admin/app/goauto/migrations" + "go-admin/app/goauto/models" + + _ "github.com/go-sql-driver/mysql" + "gorm.io/driver/mysql" + "gorm.io/gorm" + "gorm.io/gorm/clause" + "gorm.io/gorm/logger" +) + +// dsnPathReplacer swaps the database name in a Go MySQL DSN of the form +// user:pass@tcp(host:port)/dbname?params — used to connect first to the +// server (no specific throwaway database yet) and then to the freshly +// created throwaway database. +var dsnPathReplacer = regexp.MustCompile(`^(.*/)([^/?]*)(\?.*)?$`) + +func dsnWithDatabase(dsn, dbName string) string { + if dsnPathReplacer.MatchString(dsn) { + return dsnPathReplacer.ReplaceAllString(dsn, "${1}"+dbName+"${3}") + } + return dsn +} + +// setupMySQLIntegrationDB is #340 phase 4 review item 2's throwaway-database +// harness: it never touches an existing database. GOAUTO_IT_MYSQL_DSN must +// point at a MySQL SERVER (any connectable path, e.g. the system "mysql" +// database) with permission to CREATE/DROP DATABASE; a uniquely named +// zz_goauto_it_340_ database is created, migrated, and guaranteed +// dropped via t.Cleanup even if the test fails or panics. +func setupMySQLIntegrationDB(t *testing.T) (dsn string, dbName string) { + t.Helper() + baseDSN := os.Getenv("GOAUTO_IT_MYSQL_DSN") + if baseDSN == "" { + t.Skip("GOAUTO_IT_MYSQL_DSN not set; skipping MySQL concurrency integration test") + } + admin, err := sql.Open("mysql", baseDSN) + if err != nil { + t.Fatalf("open admin connection: %v", err) + } + if err := admin.Ping(); err != nil { + admin.Close() + t.Fatalf("ping MySQL server: %v", err) + } + dbName = fmt.Sprintf("zz_goauto_it_340_%d", time.Now().UnixNano()) + if _, err := admin.Exec("CREATE DATABASE `" + dbName + "`"); err != nil { + admin.Close() + t.Fatalf("create throwaway database %s: %v", dbName, err) + } + t.Cleanup(func() { + defer admin.Close() + if _, err := admin.Exec("DROP DATABASE IF EXISTS `" + dbName + "`"); err != nil { + t.Errorf("failed to drop throwaway database %s (manual cleanup required): %v", dbName, err) + } + }) + + dsn = dsnWithDatabase(baseDSN, dbName) + gdb, err := gorm.Open(mysql.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)}) + if err != nil { + t.Fatalf("open throwaway database: %v", err) + } + if err := migrations.Migrate(gdb); err != nil { + t.Fatalf("migrate throwaway database: %v", err) + } + return dsn, dbName +} + +func newMySQLIntegrationConn(t *testing.T, dsn string) *gorm.DB { + t.Helper() + conn, err := gorm.Open(mysql.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)}) + if err != nil { + t.Fatalf("open MySQL connection: %v", err) + } + return conn +} + +// TestRecomputeConcurrentPurchaseTaskUnderRealMySQL is #340 phase 4 review +// item 2: it reproduces, against a real MySQL server under REPEATABLE-READ, +// the exact race writeRecomputeChanges' locking rechecks exist to close. +// +// Timeline: +// 1. Connection A begins a transaction and runs the planning step +// (recomputeChanges) — this is A's FIRST read, so it fixes A's +// REPEATABLE-READ snapshot with zero purchase_task rows. +// 2. Connection B, concurrently, takes the SAME row's FOR UPDATE lock, +// confirmed via a channel before A is allowed to proceed. +// 3. A's write phase (writeRecomputeChanges) is started in a goroutine; it +// must BLOCK trying to take the same FOR UPDATE lock B already holds — +// the test asserts A has NOT finished after a wait window, proving a +// real block happened (not just a fast, uncontended lock grant). +// 4. B inserts a purchase_task for the row and commits, releasing the lock. +// 5. A's write phase unblocks, re-checks purchase_task under lock, and must +// see B's now-committed row and skip — this only holds because the +// recheck is a locking (FOR SHARE) read; a plain COUNT(*) would still be +// bound to A's step-1 snapshot (zero rows) and would wrongly write. +func TestRecomputeConcurrentPurchaseTaskUnderRealMySQL(t *testing.T) { + dsn, _ := setupMySQLIntegrationDB(t) + seedConn := newMySQLIntegrationConn(t, dsn) + + rule := models.SYBProductFilter{Kind: "keyword", Keyword: "档口", NormalizedKeyword: "档口", Enabled: true} + if err := seedConn.Create(&rule).Error; err != nil { + t.Fatal(err) + } + row := models.SYBProduct{ + OrderCode: "ORD-IT-TASK", DetailID: 1, StockID: 1, ShopeeItemID: "1", Quantity: 1, + ParseStatus: models.SYBParseStatusSuccess, RawJSON: `{"variationSku":"档口-1"}`, + } + if err := seedConn.Create(&row).Error; err != nil { + t.Fatal(err) + } + // purchase_task.pdd_product_id has a real FK (unlike this package's + // SQLite-backed tests, which don't enable foreign key enforcement) — + // MySQL requires an actual pdd_product row to reference. + pdd := models.PDDProduct{GoodsID: "IT-PDD-1", URL: "https://example.invalid/it", SpecsJSON: "[]"} + if err := seedConn.Create(&pdd).Error; err != nil { + t.Fatal(err) + } + + connA := newMySQLIntegrationConn(t, dsn) + connB := newMySQLIntegrationConn(t, dsn) + ctx := context.Background() + + txA := connA.Begin() + // Guard against ANY early return (t.Fatalf, panic) leaving txA open: an + // abandoned open transaction holds a connection into this throwaway + // database and blocks the DROP DATABASE cleanup indefinitely. Rollback + // on an already-committed transaction is a harmless no-op error, which + // is why the plain Commit() path below intentionally does not disable + // this cleanup. + t.Cleanup(func() { txA.Rollback() }) + _, planned, err := recomputeChanges(ctx, txA) + if err != nil { + t.Fatalf("plan: %v", err) + } + if len(planned) != 1 || planned[0].id != row.ID { + t.Fatalf("expected exactly the seeded row to be planned, got %+v", planned) + } + + txB := connB.Begin() + t.Cleanup(func() { txB.Rollback() }) + + bHoldingLock := make(chan struct{}) + bCanCommit := make(chan struct{}) + bDone := make(chan error, 1) + go func() { + var locked models.SYBProduct + if err := txB.Clauses(clause.Locking{Strength: clause.LockingStrengthUpdate}). + First(&locked, row.ID).Error; err != nil { + bDone <- fmt.Errorf("B lock row: %w", err) + return + } + close(bHoldingLock) + <-bCanCommit + task := models.PurchaseTask{ + SYBProductID: &row.ID, PDDProductID: pdd.ID, Quantity: 1, CreateRequestID: "it-race-task", + Status: models.PurchaseTaskStatusFailed, ExecutionMode: models.PurchaseExecutionModeLive, + TaskType: models.PurchaseTaskTypeSYBOrder, RuleSnapshot: "{}", + } + if err := txB.Create(&task).Error; err != nil { + bDone <- fmt.Errorf("B insert task: %w", err) + return + } + bDone <- txB.Commit().Error + }() + + select { + case <-bHoldingLock: + case err := <-bDone: + t.Fatalf("B failed before taking the row lock: %v", err) + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for B to take the row lock") + } + + aDone := make(chan struct{}) + var aActual RecomputeCounts + var aErr error + go func() { + aActual, aErr = writeRecomputeChanges(ctx, txA, planned) + close(aDone) + }() + + // A must still be blocked on B's row lock at this point — this is the + // test's proof that a real MySQL row lock, not just program logic, is + // what's being exercised. + select { + case <-aDone: + t.Fatal("A's write phase returned before B committed — it should have blocked on the row's FOR UPDATE lock") + case <-time.After(300 * time.Millisecond): + } + + close(bCanCommit) + if err := <-bDone; err != nil { + t.Fatalf("B failed: %v", err) + } + + select { + case <-aDone: + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for A's write phase to unblock after B committed") + } + if aErr != nil { + t.Fatalf("A's write phase failed: %v", aErr) + } + if err := txA.Commit().Error; err != nil { + t.Fatalf("commit A: %v", err) + } + + if aActual.SkippedHasTask != 1 || aActual.PDDToExcluded != 0 { + t.Fatalf("expected A to skip the row for the concurrently-created task, got %+v", aActual) + } + var reloaded models.SYBProduct + if err := seedConn.First(&reloaded, row.ID).Error; err != nil { + t.Fatal(err) + } + if reloaded.PDDExcluded { + t.Fatalf("row must not have been marked excluded — the concurrent task should have blocked it: %+v", reloaded) + } +} + +// TestRecomputeConcurrentReturnMatchUnderRealMySQL is the same scenario as +// TestRecomputeConcurrentPurchaseTaskUnderRealMySQL, with an active +// return_match row instead of a purchase_task as B's concurrent write. +func TestRecomputeConcurrentReturnMatchUnderRealMySQL(t *testing.T) { + dsn, _ := setupMySQLIntegrationDB(t) + seedConn := newMySQLIntegrationConn(t, dsn) + + rule := models.SYBProductFilter{Kind: "keyword", Keyword: "档口", NormalizedKeyword: "档口", Enabled: true} + if err := seedConn.Create(&rule).Error; err != nil { + t.Fatal(err) + } + row := models.SYBProduct{ + OrderCode: "ORD-IT-MATCH", DetailID: 1, StockID: 1, ShopeeItemID: "1", Quantity: 1, + ParseStatus: models.SYBParseStatusSuccess, RawJSON: `{"variationSku":"档口-1"}`, + } + if err := seedConn.Create(&row).Error; err != nil { + t.Fatal(err) + } + yeekeItem := models.YeekeReturnItem{PackageID: 0, ExternalKey: "it-race-return", ItemID: "1", VariationName: "档口-1", LastSyncedAt: time.Now()} + // A package row is required by the return_match/yeeke schema's foreign + // key; seed a minimal one. + pkg := models.YeekeReturnPackage{ExternalID: "it-race-pkg", OrderSN: "IT-ORD", TrackingNo: "IT-TRK", LastSyncedAt: time.Now()} + if err := seedConn.Create(&pkg).Error; err != nil { + t.Fatal(err) + } + yeekeItem.PackageID = pkg.ID + if err := seedConn.Create(&yeekeItem).Error; err != nil { + t.Fatal(err) + } + + connA := newMySQLIntegrationConn(t, dsn) + connB := newMySQLIntegrationConn(t, dsn) + ctx := context.Background() + + txA := connA.Begin() + // See TestRecomputeConcurrentPurchaseTaskUnderRealMySQL for why this + // unconditional cleanup is necessary regardless of the Commit() below. + t.Cleanup(func() { txA.Rollback() }) + _, planned, err := recomputeChanges(ctx, txA) + if err != nil { + t.Fatalf("plan: %v", err) + } + if len(planned) != 1 || planned[0].id != row.ID { + t.Fatalf("expected exactly the seeded row to be planned, got %+v", planned) + } + + txB := connB.Begin() + t.Cleanup(func() { txB.Rollback() }) + + bHoldingLock := make(chan struct{}) + bCanCommit := make(chan struct{}) + bDone := make(chan error, 1) + go func() { + var locked models.SYBProduct + if err := txB.Clauses(clause.Locking{Strength: clause.LockingStrengthUpdate}). + First(&locked, row.ID).Error; err != nil { + bDone <- fmt.Errorf("B lock row: %w", err) + return + } + close(bHoldingLock) + <-bCanCommit + match := models.ReturnMatch{ + SYBProductID: row.ID, YeekeReturnItemID: yeekeItem.ID, + ActiveSYBProductID: &row.ID, Status: models.ReturnMatchStatusMatched, MatchedAt: time.Now(), + } + if err := txB.Create(&match).Error; err != nil { + bDone <- fmt.Errorf("B insert match: %w", err) + return + } + bDone <- txB.Commit().Error + }() + + select { + case <-bHoldingLock: + case err := <-bDone: + t.Fatalf("B failed before taking the row lock: %v", err) + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for B to take the row lock") + } + + aDone := make(chan struct{}) + var aActual RecomputeCounts + var aErr error + go func() { + aActual, aErr = writeRecomputeChanges(ctx, txA, planned) + close(aDone) + }() + + select { + case <-aDone: + t.Fatal("A's write phase returned before B committed — it should have blocked on the row's FOR UPDATE lock") + case <-time.After(300 * time.Millisecond): + } + + close(bCanCommit) + if err := <-bDone; err != nil { + t.Fatalf("B failed: %v", err) + } + + select { + case <-aDone: + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for A's write phase to unblock after B committed") + } + if aErr != nil { + t.Fatalf("A's write phase failed: %v", aErr) + } + if err := txA.Commit().Error; err != nil { + t.Fatalf("commit A: %v", err) + } + + if aActual.SkippedReturnMatch != 1 || aActual.PDDToExcluded != 0 { + t.Fatalf("expected A to skip the row for the concurrently-created return match, got %+v", aActual) + } + var reloaded models.SYBProduct + if err := seedConn.First(&reloaded, row.ID).Error; err != nil { + t.Fatal(err) + } + if reloaded.PDDExcluded { + t.Fatalf("row must not have been marked excluded — the concurrent return match should have blocked it: %+v", reloaded) + } +} diff --git a/server/app/goauto/sybproductfilter/recompute_test.go b/server/app/goauto/sybproductfilter/recompute_test.go index bb73f5c..14018a2 100644 --- a/server/app/goauto/sybproductfilter/recompute_test.go +++ b/server/app/goauto/sybproductfilter/recompute_test.go @@ -388,3 +388,109 @@ func TestMarkedCountsByRule(t *testing.T) { t.Fatalf("expected rule B marked count 1, got %d", byID[ruleBID]) } } + +// TestRecomputeFingerprintChangesWhenRuleEvidenceChanges is #340 phase 4 +// review item 1: the fingerprint must depend on the rule's kind/keyword, not +// just its id — because those are exactly what RecomputeExecute is about to +// write into excluded_rule_kind/excluded_rule_keyword. Same rule id, same +// affected product, same direction, but the rule's own keyword changed +// between preview and execute (simulated by editing the row directly since +// the API has no edit endpoint) must be rejected as stale. +func TestRecomputeFingerprintChangesWhenRuleEvidenceChanges(t *testing.T) { + db := testDB(t) + rule := models.SYBProductFilter{Kind: "keyword", Keyword: "档口", NormalizedKeyword: "档口", Enabled: true} + if err := db.Create(&rule).Error; err != nil { + t.Fatal(err) + } + row := models.SYBProduct{OrderCode: "ORD-RULE-EDIT", DetailID: 1, StockID: 1, ShopeeItemID: "1", Quantity: 1, ParseStatus: models.SYBParseStatusSuccess, RawJSON: `{"variationSku":"档口-1"}`} + if err := db.Create(&row).Error; err != nil { + t.Fatal(err) + } + + s := NewService(db) + preview, err := s.RecomputePreview(context.Background()) + if err != nil { + t.Fatal(err) + } + if preview.PDDToExcluded != 1 || len(preview.Samples) != 1 || preview.Samples[0].RuleKeyword != "档口" { + t.Fatalf("unexpected preview: %+v", preview) + } + + // The rule's own keyword and normalized_keyword change (same id, same + // kind, still matches the same variationSku prefix) — the plan's set of + // affected ids/directions is unchanged, but the evidence that would be + // written is not. + if err := db.Model(&models.SYBProductFilter{}).Where("id = ?", rule.ID). + Updates(map[string]any{"keyword": "档口新", "normalized_keyword": "档口"}).Error; err != nil { + t.Fatal(err) + } + + _, err = s.RecomputeExecute(context.Background(), "admin1", preview.Fingerprint) + if err == nil { + t.Fatalf("expected the changed rule evidence to be rejected as stale") + } + se, ok := err.(*ServiceError) + if !ok || se.Code != CodeRecomputeStale { + t.Fatalf("expected CodeRecomputeStale, got %v", err) + } + + var reloaded models.SYBProduct + db.First(&reloaded, row.ID) + if reloaded.PDDExcluded { + t.Fatalf("nothing should have been written: %+v", reloaded) + } + var logCount int64 + db.Model(&models.SYBProductFilterRecomputeLog{}).Count(&logCount) + if logCount != 0 { + t.Fatalf("no audit log row should have been written, got %d", logCount) + } + + // A fresh preview reflects the new keyword and executes normally. + freshPreview, err := s.RecomputePreview(context.Background()) + if err != nil { + t.Fatal(err) + } + if freshPreview.Samples[0].RuleKeyword != "档口新" { + t.Fatalf("expected fresh preview to see the new keyword, got %+v", freshPreview.Samples) + } + if _, err := s.RecomputeExecute(context.Background(), "admin1", freshPreview.Fingerprint); err != nil { + t.Fatalf("fresh fingerprint should be accepted: %v", err) + } + db.First(&reloaded, row.ID) + if !reloaded.PDDExcluded || reloaded.ExcludedRuleKeyword != "档口新" { + t.Fatalf("expected the row to be excluded with the new keyword snapshot: %+v", reloaded) + } +} + +// TestRecomputeFingerprintStableAcrossUnchangedPreviews is the companion +// regression: an unchanged dataset must give the SAME fingerprint on two +// consecutive previews (map/slice iteration order must never leak into the +// hash), and that fingerprint must still execute successfully. +func TestRecomputeFingerprintStableAcrossUnchangedPreviews(t *testing.T) { + db := testDB(t) + if err := db.Create(&models.SYBProductFilter{Kind: "keyword", Keyword: "档口", NormalizedKeyword: "档口", Enabled: true}).Error; err != nil { + t.Fatal(err) + } + for i := 0; i < 5; i++ { + row := models.SYBProduct{OrderCode: fmt.Sprintf("ORD-STABLE-%d", i), DetailID: uint64(i + 1), StockID: uint64(i + 1), ShopeeItemID: fmt.Sprintf("s%d", i), Quantity: 1, ParseStatus: models.SYBParseStatusSuccess, RawJSON: `{"variationSku":"档口-x"}`} + if err := db.Create(&row).Error; err != nil { + t.Fatal(err) + } + } + + s := NewService(db) + first, err := s.RecomputePreview(context.Background()) + if err != nil { + t.Fatal(err) + } + second, err := s.RecomputePreview(context.Background()) + if err != nil { + t.Fatal(err) + } + if first.Fingerprint == "" || first.Fingerprint != second.Fingerprint { + t.Fatalf("expected a stable, non-empty fingerprint across two previews of the same data: %q vs %q", first.Fingerprint, second.Fingerprint) + } + if _, err := s.RecomputeExecute(context.Background(), "admin1", second.Fingerprint); err != nil { + t.Fatalf("unchanged-data fingerprint should execute successfully: %v", err) + } +}