feat(yeeke): sync reshipped returns and exclude new matching (#345)

This commit is contained in:
QiuSW
2026-09-28 16:30:37 +08:00
parent bce4897390
commit 3b95d4e879
13 changed files with 355 additions and 98 deletions
@@ -0,0 +1,38 @@
package returnmatch
import (
"context"
"errors"
"testing"
"time"
"go-admin/app/goauto/models"
)
func TestReshippedReturnExcludedAndRecheckedBeforeInsert(t *testing.T) {
db := testDB(t)
s := NewService(db)
deadline := time.Now().Add(24 * time.Hour)
ret := seedReturn(t, db, "红色", &deadline)
sy := seedSYB(t, db, "TEST", 1, "红色", "", time.Now())
pool, err := s.availableReturnPool(context.Background())
if err != nil || len(pool) != 1 {
t.Fatalf("waiting pool=%v err=%v", pool, err)
}
if err := db.Model(&models.YeekeReturnPackage{}).Where("id = ?", ret.PackageID).Update("claim_status", "2").Error; err != nil {
t.Fatal(err)
}
pool, err = s.availableReturnPool(context.Background())
if err != nil || len(pool) != 0 {
t.Fatalf("reshipped pool=%v err=%v", pool, err)
}
_, err = s.matchOneWithLock(context.Background(), sy.ID, MatchOutcome{ReturnItemID: ret.ID, DestroyDeadline: deadline}, "test")
if !errors.Is(err, errReturnNoLongerEligible) {
t.Fatalf("final recheck=%v", err)
}
var count int64
db.Model(&models.ReturnMatch{}).Count(&count)
if count != 0 {
t.Fatal("reshipped item was allocated")
}
}
+18
View File
@@ -238,6 +238,11 @@ func (s *Service) batchMatch(ctx context.Context, req BatchMatchRequest) (BatchM
}
match, insertErr := s.matchOneWithLock(ctx, id, outcome, req.Operator)
if insertErr != nil {
if errors.Is(insertErr, errReturnNoLongerEligible) {
resp.Items = append(resp.Items, BatchMatchItem{SYBProductID: id, ReasonCode: ReasonNoCandidate, Reason: "退货商品已重出或不再可用"})
resp.SkippedCount++
continue
}
if errors.Is(insertErr, errStageNoLongerEligible) {
// #338 review fix: the stage was re-checked under the same
// FOR UPDATE lock purchase.create takes, right before insert.
@@ -292,6 +297,7 @@ func (s *Service) availableReturnPool(ctx context.Context) ([]ReturnCandidate, e
Joins("JOIN yeeke_return_package AS p ON p.id = i.package_id").
Joins("LEFT JOIN return_match AS m ON m.active_yeeke_return_item_id = i.id").
Where("m.id IS NULL AND i.sync_status = ? AND p.sync_status = ?", "ok", "ok").
Where("p.claim_status = ? AND p.status_unrecognized = ?", "1", false).
Find(&rows).Error
if err != nil {
return nil, err
@@ -312,6 +318,7 @@ func (s *Service) availableReturnPool(ctx context.Context) ([]ReturnCandidate, e
// it between BatchMatch's outer screening pass and this point (#338 review
// fix: race between matching and purchase creation).
var errStageNoLongerEligible = errors.New("syb product stage no longer participates in matching")
var errReturnNoLongerEligible = errors.New("yeeke return is no longer waiting to ship")
// matchOneWithLock takes the SAME row lock purchase.Service.create takes on
// syb_product (clause.Locking{Strength: "UPDATE"}) and re-computes the
@@ -340,6 +347,17 @@ func (s *Service) matchOneWithLock(ctx context.Context, sybID uint64, outcome Ma
if err := tx.First(&returnItem, outcome.ReturnItemID).Error; err != nil {
return err
}
var pkg models.YeekeReturnPackage
if err := tx.Clauses(clauseLockUpdate()).First(&pkg, returnItem.PackageID).Error; err != nil {
return err
}
// Same package -> item lock order as sync upsert; avoid a lock cycle.
if err := tx.Clauses(clauseLockUpdate()).First(&returnItem, outcome.ReturnItemID).Error; err != nil {
return err
}
if pkg.ClaimStatus != "1" || pkg.StatusUnrecognized || pkg.SyncStatus != "ok" || returnItem.SyncStatus != "ok" {
return errReturnNoLongerEligible
}
sybIDCopy := syb.ID
returnIDCopy := outcome.ReturnItemID
deadline := outcome.DestroyDeadline
@@ -69,7 +69,8 @@ func seedSYB(t *testing.T, db *gorm.DB, orderCode string, detailID uint64, color
func seedReturn(t *testing.T, db *gorm.DB, variationName string, deadline *time.Time) models.YeekeReturnItem {
t.Helper()
pkg := models.YeekeReturnPackage{
ExternalID: "pkg-" + variationName + fmt.Sprint(time.Now().UnixNano()), OrderSN: "ORD1", TrackingNo: "TRK1",
ClaimStatus: "1",
ExternalID: "pkg-" + variationName + fmt.Sprint(time.Now().UnixNano()), OrderSN: "ORD1", TrackingNo: "TRK1",
DestroyDeadLine: deadline, LastSyncedAt: time.Now(),
}
if err := db.Create(&pkg).Error; err != nil {
+112
View File
@@ -0,0 +1,112 @@
package yeeke
import (
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"testing"
"go-admin/app/goauto/models"
"go-admin/app/goauto/yeekeclient"
)
func TestSyncBothStatusesPreservesIdentityAndAvailability(t *testing.T) {
db := testDB(t)
phase := 0
var calls []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var req struct {
Status string `json:"status"`
Page int `json:"pageNo"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
t.Error(err)
return
}
calls = append(calls, fmt.Sprintf("%s/%d", req.Status, req.Page))
w.Header().Set("Content-Type", "application/json")
if phase == 0 && req.Status == "2" {
fmt.Fprint(w, page(nil, 0, 0))
return
}
if req.Status == "1" {
if req.Page == 1 {
fmt.Fprint(w, page([]string{record("p1", "i1", "v1", 1)}, 2, 2))
} else {
fmt.Fprint(w, page([]string{record("p2", "i2", "v2", 1)}, 2, 2))
}
} else {
fmt.Fprint(w, page([]string{record("p1", "i1", "v1", 2)}, 1, 1))
}
}))
defer srv.Close()
c, _ := yeekeclient.New(srv.URL)
s := NewService(db, c, Config{PageSize: 1})
if _, err := s.Sync(context.Background(), "manual"); err != nil {
t.Fatal(err)
}
var original models.YeekeReturnPackage
db.Where("external_id = ?", "p1").First(&original)
phase = 1
calls = nil
for run := 0; run < 2; run++ {
rep, err := s.Sync(context.Background(), "manual")
if err != nil || rep.Status != "succeeded" || rep.MissingMarked != 0 {
t.Fatalf("rep=%+v err=%v", rep, err)
}
}
if fmt.Sprint(calls) != "[1/1 1/2 2/1 1/1 1/2 2/1]" {
t.Fatalf("independent pagination: %v", calls)
}
var current models.YeekeReturnPackage
db.First(&current, original.ID)
if current.ClaimStatus != "2" || current.StatusUnrecognized || current.SyncStatus != "ok" {
t.Fatalf("current=%+v", current)
}
var count int64
db.Model(&models.YeekeReturnPackage{}).Count(&count)
if count != 2 {
t.Fatalf("packages=%d", count)
}
db.Model(&models.YeekeReturnItem{}).Count(&count)
if count != 2 {
t.Fatalf("items=%d", count)
}
}
func TestReshipPageFailureDoesNotMarkMissing(t *testing.T) {
db := testDB(t)
old := models.YeekeReturnPackage{ExternalID: "old", ClaimStatus: "2", SyncStatus: "ok"}
if err := db.Create(&old).Error; err != nil {
t.Fatal(err)
}
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var req struct {
Status string `json:"status"`
}
json.NewDecoder(r.Body).Decode(&req)
if req.Status == "2" {
w.WriteHeader(http.StatusBadGateway)
return
}
fmt.Fprint(w, page([]string{record("new", "i", "v", 1)}, 1, 1))
}))
defer srv.Close()
c, _ := yeekeclient.New(srv.URL)
s := NewService(db, c, Config{PageSize: 10})
rep, err := s.Sync(context.Background(), "manual")
if err == nil || rep.Status != "failed" || rep.Created != 1 || rep.MissingMarked != 0 {
t.Fatalf("rep=%+v err=%v", rep, err)
}
db.First(&old, old.ID)
if old.SyncStatus != "ok" {
t.Fatal("incomplete combined sync marked reshipped package missing")
}
var run models.YeekeSyncRun
db.First(&run, rep.RunID)
if run.ErrorMessage == "" {
t.Fatal("missing state/page failure diagnostic")
}
}
+79 -64
View File
@@ -54,10 +54,8 @@ type Report struct {
}
// knownClaimStatuses lists the status values the sync code currently
// understands. The list surface (POST .../relation/list) is queried with
// status=1, so "1" is the only value observed in practice; anything else is
// flagged rather than silently accepted or rejected (#336).
var knownClaimStatuses = map[string]bool{"1": true}
// understands: waiting to ship (1) and reshipped (2), confirmed by HAR.
var knownClaimStatuses = map[string]bool{"1": true, "2": true}
func external(v any) string { return fmt.Sprint(v) }
func stamp(t *yeekeclient.Timestamp) *time.Time {
@@ -194,86 +192,103 @@ func (s *Service) run(ctx context.Context, r *models.YeekeSyncRun) (Report, erro
}
s.db.Model(r).Updates(updates)
}()
seen := map[string]bool{}
var firstWriteErr error
seenPackages := map[string]string{}
// complete tracks whether the page walk ended NATURALLY (empty page,
// short page, or reaching p.Pages) as opposed to the duplicate-
// fingerprint break or MaxPages exhaustion (#338): only a naturally
// complete run is trusted to mark absent items/packages "missing" below,
// since a duplicate/MaxPages stop means the walk never actually finished
// seeing everything yeeke currently has.
complete := false
for page := 1; page <= s.cfg.MaxPages; page++ {
var p yeekeclient.ReturnPage
var e error
for a := 0; ; a++ {
p, e = s.client.List(ctx, page, s.cfg.PageSize)
if e == nil || a >= s.cfg.Retry {
complete := true
for _, status := range []string{"1", "2"} {
seen := map[string]bool{}
statusComplete := false
for page := 1; page <= s.cfg.MaxPages; page++ {
var p yeekeclient.ReturnPage
var e error
for a := 0; ; a++ {
p, e = s.client.ListStatus(ctx, page, s.cfg.PageSize, status)
if e == nil || a >= s.cfg.Retry {
break
}
select {
case <-ctx.Done():
runErr = ctx.Err()
errMsg = truncateRunError(runErr.Error())
return rep, runErr
case <-time.After(time.Duration(a+1) * 100 * time.Millisecond):
}
}
if e != nil {
// A failed page never overwrites what earlier pages already wrote
// (#336): the run simply stops here and everything upserted so far
// stays as-is, reported through Read/Created/Updated above.
runErr = e
errMsg = truncateRunError(fmt.Sprintf("状态 %s 第 %d 页拉取失败:%v", status, page, e))
return rep, runErr
}
rep.TotalPages++
if len(p.Records) == 0 {
statusComplete = true
break
}
select {
case <-ctx.Done():
runErr = ctx.Err()
errMsg = truncateRunError(runErr.Error())
return rep, runErr
case <-time.After(time.Duration(a+1) * 100 * time.Millisecond):
finger := pageFingerprint(p)
if seen[finger] {
rep.Skipped += len(p.Records)
break
}
}
if e != nil {
// A failed page never overwrites what earlier pages already wrote
// (#336): the run simply stops here and everything upserted so far
// stays as-is, reported through Read/Created/Updated above.
runErr = e
errMsg = truncateRunError(e.Error())
return rep, runErr
}
rep.TotalPages = page
if len(p.Records) == 0 {
complete = true
break
}
finger := pageFingerprint(p)
if seen[finger] {
rep.Skipped += len(p.Records)
break
}
seen[finger] = true
for _, x := range p.Records {
created, updated, recovered, err := s.upsert(ctx, x)
if err != nil {
rep.Failed++
if firstWriteErr == nil {
firstWriteErr = err
seen[finger] = true
for _, x := range p.Records {
key, currentStatus := packageKey(x), external(x.Status)
if previous, ok := seenPackages[key]; ok && (previous == "2" || previous == currentStatus) {
rep.Skipped++
continue
}
created, updated, recovered, err := s.upsert(ctx, x)
seenPackages[key] = currentStatus
if err != nil {
rep.Failed++
if firstWriteErr == nil {
firstWriteErr = err
}
continue
}
rep.Read++
rep.Recovered += recovered
if created {
rep.Created++
} else if updated {
rep.Updated++
} else {
rep.Skipped++
}
continue
}
rep.Read++
rep.Recovered += recovered
if created {
rep.Created++
} else if updated {
rep.Updated++
} else {
rep.Skipped++
if len(p.Records) < s.cfg.PageSize {
statusComplete = true
break
}
if p.Pages > 0 && page >= p.Pages {
statusComplete = true
break
}
}
if len(p.Records) < s.cfg.PageSize {
complete = true
break
}
if p.Pages > 0 && page >= p.Pages {
complete = true
if !statusComplete {
complete = false
errMsg = fmt.Sprintf("状态 %s 分页未完整结束(重复页或达到页数上限),未执行缺失标记", status)
break
}
}
if !complete {
runErr = errors.New(errMsg)
return rep, runErr
}
rep.Status = "succeeded"
if rep.Failed > 0 && firstWriteErr != nil {
// Surface why records failed instead of a bare counter.
errMsg = truncateRunError(fmt.Sprintf("%d 条写入失败,首个原因:%v", rep.Failed, firstWriteErr))
if rep.Read == 0 {
rep.Status = "failed"
runErr = errors.New(errMsg)
}
rep.Status = "failed"
runErr = errors.New(errMsg)
}
// #338: only a naturally complete run with zero write failures is
// trusted to mark items/packages the sync no longer sees as "missing".
@@ -288,7 +303,7 @@ func (s *Service) run(ctx context.Context, r *models.YeekeSyncRun) (Report, erro
rep.MissingMarked = marked
}
}
return rep, nil
return rep, runErr
}
// markMissing implements #338's completion-triggered availability flip: any
+4 -10
View File
@@ -204,11 +204,8 @@ func TestDuplicateFingerprintStopsMarking(t *testing.T) {
defer srv2.Close()
s.client, _ = yeekeclient.New(srv2.URL)
rep, err := s.Sync(context.Background(), "manual")
if err != nil {
t.Fatalf("second sync: %v", err)
}
if rep.MissingMarked != 0 {
t.Fatalf("MissingMarked=%d, want 0 (duplicate-fingerprint stop is not complete)", rep.MissingMarked)
if err == nil || rep.Status != "failed" || rep.MissingMarked != 0 {
t.Fatalf("rep=%+v err=%v (duplicate-fingerprint stop is not complete)", rep, err)
}
var p2 models.YeekeReturnPackage
@@ -250,11 +247,8 @@ func TestMaxPagesExhaustionStopsMarking(t *testing.T) {
s.client, _ = yeekeclient.New(srv2.URL)
s.cfg.MaxPages = 2
rep, err := s.Sync(context.Background(), "manual")
if err != nil {
t.Fatalf("second sync: %v", err)
}
if rep.MissingMarked != 0 {
t.Fatalf("MissingMarked=%d, want 0 (MaxPages exhaustion is not complete)", rep.MissingMarked)
if err == nil || rep.Status != "failed" || rep.MissingMarked != 0 {
t.Fatalf("rep=%+v err=%v (MaxPages exhaustion is not complete)", rep, err)
}
var p2 models.YeekeReturnPackage
+2 -5
View File
@@ -108,10 +108,7 @@ func TestPagingSkipsARepeatedDuplicatePage(t *testing.T) {
c, _ := yeekeclient.New(srv.URL)
s := NewService(db, c, Config{PageSize: 1})
rep, err := s.Sync(context.Background(), "manual")
if err != nil {
t.Fatal(err)
}
if rep.Status != "succeeded" {
if err == nil || rep.Status != "failed" {
t.Fatalf("rep=%+v", rep)
}
var n int64
@@ -138,7 +135,7 @@ func TestPagingStopsOnEmptyPage(t *testing.T) {
if err != nil {
t.Fatal(err)
}
if rep.Status != "succeeded" || rep.TotalPages != 1 || rep.Read != 0 {
if rep.Status != "succeeded" || rep.TotalPages != 2 || rep.Read != 0 {
t.Fatalf("rep=%+v", rep)
}
}
+9 -1
View File
@@ -371,8 +371,16 @@ func (f *FlexInt) UnmarshalJSON(b []byte) error {
}
func (c *Client) List(ctx context.Context, pageNo, pageSize int) (ReturnPage, error) {
return c.ListStatus(ctx, pageNo, pageSize, "1")
}
// ListStatus reads only the two HAR-confirmed return statuses.
func (c *Client) ListStatus(ctx context.Context, pageNo, pageSize int, status string) (ReturnPage, error) {
if status != "1" && status != "2" {
return ReturnPage{}, fmt.Errorf("unsupported yeeke return status")
}
// Same shape the web client posts (HAR): sort via column/order, filters as strings.
body := map[string]any{"pageNo": pageNo, "pageSize": pageSize, "claimFlag": "1", "status": "1", "relationFlag": "1", "column": "createTime", "order": "desc"}
body := map[string]any{"pageNo": pageNo, "pageSize": pageSize, "claimFlag": "1", "status": status, "relationFlag": "1", "column": "createTime", "order": "desc"}
raw, e := c.do(ctx, http.MethodPost, "/agent-foreign/packageClaimRec/relation/list", body, nil)
if e != nil {
return ReturnPage{}, e
@@ -0,0 +1,36 @@
package yeekeclient
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
)
func TestListStatusUsesConfirmedHARFilters(t *testing.T) {
var statuses []string
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var body map[string]any
json.NewDecoder(r.Body).Decode(&body)
if r.Method != http.MethodPost || r.URL.Path != "/agent-foreign/packageClaimRec/relation/list" || body["claimFlag"] != "1" || body["relationFlag"] != "1" || body["column"] != "createTime" || body["order"] != "desc" {
t.Error("HAR filters changed")
}
statuses = append(statuses, body["status"].(string))
w.Header().Set("Content-Type", "application/json")
w.Write([]byte(`{"success":true,"result":{"records":[],"pages":0,"total":0}}`))
}))
defer s.Close()
c, _ := New(s.URL)
for _, status := range []string{"1", "2"} {
if _, err := c.ListStatus(context.Background(), 1, 20, status); err != nil {
t.Fatal(err)
}
}
if _, err := c.ListStatus(context.Background(), 1, 20, "3"); err == nil {
t.Fatal("unsupported status allowed")
}
if len(statuses) != 2 || statuses[0] != "1" || statuses[1] != "2" {
t.Fatal(statuses)
}
}