Merge remote-tracking branch 'origin/main' into feat/340-syb-excluded-products

# Conflicts:
#	server/app/goauto/sybimport/handler.go
#	server/app/goauto/sybimport/service.go
#	server/app/goauto/sybimport/service_test.go
#	web/src/views/goauto/syb-products/index.vue
This commit is contained in:
QiuSW
2026-09-29 10:50:54 +08:00
61 changed files with 2907 additions and 274 deletions
+34 -6
View File
@@ -25,12 +25,8 @@ const SPADirEnv = "GOAUTO_WEB_DIST"
// dist; a NoRoute handler installed anyway would turn every genuine 404 into
// an HTML page, which is far more confusing than a plain 404.
func InitSPARouter(engine *gin.Engine) {
dist := strings.TrimSpace(os.Getenv(SPADirEnv))
if dist == "" {
dist = "dist"
}
index := filepath.Join(dist, "index.html")
if _, err := os.Stat(index); err != nil {
dist, index, ok := spaIndex()
if !ok {
return
}
@@ -56,6 +52,38 @@ func InitSPARouter(engine *gin.Engine) {
})
}
// spaIndex resolves the built frontend directory and reports whether its
// index.html exists.
func spaIndex() (dist, index string, ok bool) {
dist = strings.TrimSpace(os.Getenv(SPADirEnv))
if dist == "" {
dist = "dist"
}
index = filepath.Join(dist, "index.html")
if _, err := os.Stat(index); err != nil {
return dist, index, false
}
return dist, index, true
}
// registerRootRoute decides what `GET /` returns (#346).
//
// `[必须]` When the built frontend exists, `/` must be the Admin SPA. go-admin's
// welcome page used to own `/` in every non-prod mode, so any reverse proxy
// that forwarded `/` to this server (instead of serving dist itself) showed
// "GO-ADMIN欢迎您" instead of the Admin — which is exactly what happened after
// the 2026-09-28 server migration. The welcome page is kept only for
// development without a dist, where the frontend runs under vite.
func registerRootRoute(r gin.IRoutes, mode string, welcome gin.HandlerFunc) {
if _, index, ok := spaIndex(); ok {
r.GET("/", func(c *gin.Context) { c.File(index) })
return
}
if mode != "prod" {
r.GET("/", welcome)
}
}
// isAPIPath reports whether a path belongs to the server rather than the SPA.
func isAPIPath(path string) bool {
for _, prefix := range []string{"/api/", "/swagger/", "/static/", "/form-generator/", "/ws/", "/wslogout/", "/info"} {
+43
View File
@@ -91,3 +91,46 @@ func TestWithoutDistNoFallbackIsInstalled(t *testing.T) {
t.Fatalf("没有 dist 时接口仍应正常: %d", response.Code)
}
}
// #346: with a built frontend, `/` must be the Admin SPA — never go-admin's
// welcome page, even in non-prod modes where the welcome page used to own `/`.
func TestRootServesSPAWhenDistExists(t *testing.T) {
gin.SetMode(gin.TestMode)
dist := filepath.Join(t.TempDir(), "dist")
if err := os.MkdirAll(dist, 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(dist, "index.html"), []byte("<!doctype html>SPA"), 0o644); err != nil {
t.Fatal(err)
}
t.Setenv(SPADirEnv, dist)
for _, mode := range []string{"dev", "test", "prod"} {
engine := gin.New()
registerRootRoute(engine, mode, func(c *gin.Context) { c.String(http.StatusOK, "GO-ADMIN欢迎您") })
InitSPARouter(engine)
response := do(engine, http.MethodGet, "/")
if response.Code != http.StatusOK || response.Body.String() != "<!doctype html>SPA" {
t.Fatalf("mode=%s: / should serve index.html, got %d %q", mode, response.Code, response.Body.String())
}
}
}
// Without a dist (development under vite) the previous behaviour is kept:
// welcome page outside prod, nothing registered in prod.
func TestRootWithoutDistKeepsPreviousBehaviour(t *testing.T) {
gin.SetMode(gin.TestMode)
t.Setenv(SPADirEnv, filepath.Join(t.TempDir(), "missing-dist"))
welcome := func(c *gin.Context) { c.String(http.StatusOK, "GO-ADMIN欢迎您") }
dev := gin.New()
registerRootRoute(dev, "dev", welcome)
if response := do(dev, http.MethodGet, "/"); response.Code != http.StatusOK || response.Body.String() != "GO-ADMIN欢迎您" {
t.Fatalf("dev without dist should keep the welcome page, got %d %q", response.Code, response.Body.String())
}
prod := gin.New()
registerRootRoute(prod, "prod", welcome)
if response := do(prod, http.MethodGet, "/"); response.Code != http.StatusNotFound {
t.Fatalf("prod without dist should not register /, got %d", response.Code)
}
}
+1 -3
View File
@@ -40,9 +40,7 @@ func sysBaseRouter(r *gin.RouterGroup) {
go ws.WebsocketManager.SendService()
go ws.WebsocketManager.SendAllService()
if config.ApplicationConfig.Mode != "prod" {
r.GET("/", apis.GoAdmin)
}
registerRootRoute(r, config.ApplicationConfig.Mode, apis.GoAdmin)
r.GET("/info", handler.Ping)
}
+2
View File
@@ -173,6 +173,8 @@ func moduleKeyForAPI(path string) string {
return ModuleYeekeSyncRuns
case strings.HasPrefix(path, "/api/admin/v1/yeeke-returns"):
return ModuleYeekeReturns
case strings.HasPrefix(path, "/api/admin/v1/return-matches"):
return ModuleSYBProducts
default:
return ""
}
+12
View File
@@ -3,6 +3,9 @@ package access
// RolePurchaser is the fixed role key used by the GoAuto purchaser account.
const RolePurchaser = "purchaser"
// RoleAfterSales is the fixed role key for the internal after-sales users.
const RoleAfterSales = "after_sales"
// APIPermission describes one admin API known to GoAuto. Purchaser marks the
// APIs that the purchaser role may call; every other API remains admin-only.
type APIPermission struct {
@@ -132,6 +135,15 @@ var AdminAPIs = []APIPermission{
{"查看 yeeke 同步详情", "/api/admin/v1/yeeke-returns/sync-runs/:runId", "GET", true},
{"手动触发 yeeke 同步", "/api/admin/v1/yeeke-returns/sync", "POST", true},
{"查看退货匹配", "/api/admin/v1/return-matches", "GET", true},
{"查看退货匹配批次", "/api/admin/v1/return-matches/batches", "GET", true},
{"查看退货匹配批次详情", "/api/admin/v1/return-matches/batches/:batchId", "GET", true},
{"查看退货匹配详情", "/api/admin/v1/return-matches/:id", "GET", true},
{"批量匹配退货", "/api/admin/v1/return-matches/batch-match", "POST", true},
{"确认退货匹配", "/api/admin/v1/return-matches/:id/confirm", "POST", true},
{"取消退货匹配", "/api/admin/v1/return-matches/:id/cancel", "POST", true},
{"备注退货匹配", "/api/admin/v1/return-matches/:id/remark", "POST", true},
{"查看 AI 匹配状态", "/api/admin/v1/ai-matching-settings", "GET", true},
{"保存 AI 匹配设置", "/api/admin/v1/ai-matching-settings", "PUT", false},
{"测试 AI 服务连接", "/api/admin/v1/ai-matching-settings/test", "POST", false},
+1
View File
@@ -41,6 +41,7 @@ func MigratedModels() []any {
&models.SYBSpecAIParseRun{},
&models.SYBSpecAIParseWorkItem{},
&models.SYBSession{},
&models.SYBSessionAuthLease{},
&models.SYBShop{},
&models.SYBProductFilter{},
&models.SYBProductFilterRecomputeLog{},
+11
View File
@@ -504,6 +504,17 @@ type SYBSession struct {
UpdatedAt time.Time `json:"updatedAt"`
}
// SYBSessionAuthLease serializes re-authentication across sync and writeback
// workers. Cookies and passwords never live in this table.
type SYBSessionAuthLease struct {
ID uint64 `gorm:"primaryKey"`
Username string `gorm:"size:128;not null;uniqueIndex"`
Owner string `gorm:"size:36;not null"`
ExpiresAt time.Time `gorm:"not null"`
}
func (SYBSessionAuthLease) TableName() string { return "syb_session_auth_lease" }
// SYBProduct is one SYB (顺云宝 ERP) shipment detail line: one order can carry
// several Shopee product lines, and the same Shopee product can appear more
// than once within one order at different colors/sizes/quantities — each such
+2 -2
View File
@@ -475,10 +475,10 @@ func allowedOperator(c *gin.Context) bool {
return true
}
role, _ := jwt.ExtractClaims(c)["rolekey"].(string)
if role == "admin" || role == "purchaser" {
if role == "admin" || role == "purchaser" || role == "after_sales" {
return true
}
c.JSON(http.StatusForbidden, gin.H{"code": "FORBIDDEN", "message": "只有管理员或采购员可以操作采购任务"})
c.JSON(http.StatusForbidden, gin.H{"code": "FORBIDDEN", "message": "只有管理员、采购员或售后可以操作采购任务"})
c.Abort()
return false
}
@@ -0,0 +1,26 @@
package purchase
import (
"net/http/httptest"
"testing"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
)
func TestAllowedOperatorIncludesAfterSales(t *testing.T) {
for _, role := range []string{"admin", "purchaser", "after_sales", "other", ""} {
t.Run(role, func(t *testing.T) {
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Set("JWT_PAYLOAD", jwt.MapClaims{"rolekey": role})
want := role == "admin" || role == "purchaser" || role == "after_sales"
if got := allowedOperator(c); got != want {
t.Fatalf("role %q allowed=%v want %v", role, got, want)
}
if !want && w.Code != 403 {
t.Fatalf("denied role status=%d", w.Code)
}
})
}
}
@@ -96,7 +96,7 @@ func TestOrderBackfillMixedBatchAndReplay(t *testing.T) {
if saved.StatusVersion != a.StatusVersion+1 || saved.ErrorCode != nil || saved.ErrorMessage != nil || saved.DeviceRunSlot != nil || saved.AccountRunSlot != nil || saved.ActiveSlot == nil || saved.Status != models.PurchaseTaskStatusOrderCreated {
t.Fatalf("state metadata: %+v", saved)
}
if saved.PaymentReviewStatus != a.PaymentReviewStatus || saved.LogisticsStatus != a.LogisticsStatus || saved.WritebackStatus != a.WritebackStatus || saved.RuleSnapshot != a.RuleSnapshot {
if saved.PaymentReviewStatus != a.PaymentReviewStatus || saved.LogisticsStatus != a.LogisticsStatus || saved.WritebackStatus != models.PurchaseWritebackStatusPending || saved.RuleSnapshot != a.RuleSnapshot {
t.Fatal("unrelated business facts changed")
}
for _, replayID := range []string{rid, uuid.NewString()} {
@@ -60,7 +60,10 @@ func ensureOrderWriteback(tx *gorm.DB, t models.PurchaseTask) error {
return nil
}
row := models.PurchaseOrderWriteback{PurchaseTaskID: t.ID, StockID: int64(syb.StockID), DetailID: int64(syb.DetailID), OrderNo: *t.PDDOrderNo, Status: "pending"}
return tx.Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent)}).Clauses(clause.OnConflict{DoNothing: true}).Create(&row).Error
if err := tx.Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent)}).Clauses(clause.OnConflict{DoNothing: true}).Create(&row).Error; err != nil {
return err
}
return tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ? AND writeback_status = ?", t.ID, models.PurchaseWritebackStatusNotSelected).Update("writeback_status", models.PurchaseWritebackStatusPending).Error
}
func (s *Service) OrderWritebackViews(ctx context.Context, tasks []models.PurchaseTask) (map[uint64]OrderWritebackView, error) {
@@ -116,7 +116,11 @@ func TestOrderWritebackRemoteOutcomes(t *testing.T) {
t.Fatal("automatically repeated write")
}
after := loadBackfillTask(t, s.DB, task.ID)
if after.PaymentReviewStatus != task.PaymentReviewStatus || after.WritebackStatus != task.WritebackStatus || after.StatusVersion != task.StatusVersion {
wantTaskWriteback := models.PurchaseWritebackStatusFailed
if tc.want == "succeeded" {
wantTaskWriteback = models.PurchaseWritebackStatusSucceeded
}
if after.PaymentReviewStatus != task.PaymentReviewStatus || after.WritebackStatus != wantTaskWriteback || after.StatusVersion != task.StatusVersion {
t.Fatal("changed purchase/payment/logistics facts")
}
})
@@ -84,38 +84,16 @@ func sessionUnavailableMessage(err error) string {
return "SYB会话不可用(" + category + "),将自动重试;如持续失败请恢复登录后重试"
}
// restoreOrderWritebackClient rebuilds a SYB client from the cached session
// only. It never logs in, never triggers OCR and never deletes the cached
// session (that stays the exclusive responsibility of sybimport.Connect's
// login/refresh path) — it only reports whether the cached cookies still
// work, via CheckSession, so the caller can classify the failure (#330).
// restoreOrderWritebackClient uses the same session acquisition path as SYB
// sync. A valid cached session is reused; an explicitly invalid session is
// refreshed under the shared database auth lease so concurrent workers do not
// request multiple captcha codes.
func restoreOrderWritebackClient(ctx context.Context, db *gorm.DB) (OrderNumberClient, error) {
cfg := config.ExtConfig.SYB.Resolved()
session, err := sybclient.NewSessionStore(db).Load(ctx, cfg.Username, time.Now())
if err != nil {
return nil, err
}
if session.UserID <= 0 {
return nil, errSessionUserIDMissing
}
c, err := sybclient.New(cfg.BaseURL)
if err != nil {
return nil, err
}
if err = c.ImportCookiesJSON(session.CookiesJSON); err != nil {
return nil, err
}
// Active probe (#330 修订1): without this, a remotely-expired cookie jar
// imports cleanly and only fails later inside read(), which would record
// it as SYB_READ_FAILED instead of the retryable session-class outcome.
// Any error here — ErrSessionInvalid or network/format — is treated as
// session-class; only ErrSessionInvalid is a confirmed logout, but a
// network/format error is not confirmed-valid either, so it is still
// retried rather than attempted as a write.
if err = c.CheckSession(ctx, session.UserID, cfg.Username); err != nil {
return nil, err
}
return c, nil
return sybclient.AcquireSession(ctx, db, sybclient.LoginConfig{
BaseURL: cfg.BaseURL, Username: cfg.Username, Password: cfg.Password,
OcrURL: cfg.OcrURL, OcrMaxAttempts: cfg.OcrMaxAttempts,
})
}
// One short-lived claim at a time across processes. No business writes occur
@@ -184,7 +162,20 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
if status == "succeeded" {
updates["completed_at"] = w.Now()
}
return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error
if err := db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error; err != nil {
return err
}
// Keep the admin task state aligned with the authoritative writeback
// record. Conflicts and unknown outcomes remain failed until a human
// resolves them; they must never appear as successful.
taskStatus := models.PurchaseWritebackStatusFailed
if status == "succeeded" {
taskStatus = models.PurchaseWritebackStatusSucceeded
}
return db.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", item.PurchaseTaskID).Updates(map[string]any{
"writeback_status": taskStatus,
"writeback_at": gorm.Expr("CASE WHEN ? = 'succeeded' THEN ? ELSE writeback_at END", status, w.Now()),
}).Error
}
// finishSessionUnavailable is the bounded-retry counterpart of finish for
// SYB_SESSION_UNAVAILABLE: instead of clearing the lease, it schedules the
@@ -201,7 +192,10 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
} else {
updates["lease_expires_at"] = nil
}
return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error
if err := db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error; err != nil {
return err
}
return db.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", item.PurchaseTaskID).Update("writeback_status", models.PurchaseWritebackStatusFailed).Error
}
var task models.PurchaseTask
if err = db.First(&task, item.PurchaseTaskID).Error; err != nil {
@@ -0,0 +1,120 @@
package returnmatch
import (
"context"
"errors"
"fmt"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"go-admin/app/goauto/models"
"gorm.io/gorm"
)
func seedCancelMatch(t *testing.T, db *gorm.DB, status string) models.ReturnMatch {
t.Helper()
syb := seedSYB(t, db, "CANCEL-TEST", 1, "白色", "L", time.Now())
item := seedReturn(t, db, "白色,L", nil)
m := models.ReturnMatch{SYBProductID: syb.ID, YeekeReturnItemID: item.ID, ActiveSYBProductID: &syb.ID, ActiveYeekeReturnItemID: &item.ID, Status: status, MatchedAt: time.Now()}
if err := db.Create(&m).Error; err != nil {
t.Fatal(err)
}
return m
}
func TestCancelExpectedStatusHTTP(t *testing.T) {
for _, tc := range []struct {
name, status, body, role string
code int
}{
{"matched", "matched", `{"expectedStatus":"matched"}`, "purchaser", 200},
{"confirmed_conflict", "confirmed", `{"expectedStatus":"matched"}`, "after_sales", 409},
{"cancelled_conflict", "cancelled", `{"expectedStatus":"matched"}`, "admin", 409},
{"legacy_empty", "confirmed", "", "after_sales", 200},
{"legacy_object", "confirmed", `{}`, "admin", 200},
{"invalid_status", "matched", `{"expectedStatus":"confirmed"}`, "admin", 400},
{"invalid_empty", "matched", `{"expectedStatus":""}`, "admin", 400},
{"invalid_null", "matched", `{"expectedStatus":null}`, "admin", 400},
{"malformed", "matched", `{`, "admin", 400},
{"forbidden", "matched", `{"expectedStatus":"matched"}`, "other", 403},
} {
t.Run(tc.name, func(t *testing.T) {
db := testDB(t)
m := seedCancelMatch(t, db, tc.status)
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Request = httptest.NewRequest("POST", fmt.Sprintf("/return-matches/%d/cancel", m.ID), strings.NewReader(tc.body))
c.Params = gin.Params{{Key: "id", Value: fmt.Sprint(m.ID)}}
c.Set("JWT_PAYLOAD", jwt.MapClaims{"rolekey": tc.role, "username": "cancel-tester"})
Handler{DB: db}.Cancel(c)
if w.Code != tc.code {
t.Fatalf("status=%d body=%s", w.Code, w.Body.String())
}
var saved models.ReturnMatch
db.First(&saved, m.ID)
var logs int64
db.Model(&models.ReturnMatchLog{}).Where("match_id = ?", m.ID).Count(&logs)
if tc.code == 200 {
if saved.Status != "cancelled" || saved.ActiveSYBProductID != nil || saved.ActiveYeekeReturnItemID != nil || saved.CancelledBy != "cancel-tester" || saved.CancelledAt == nil || logs != 1 {
t.Fatalf("cancel or audit incomplete: %+v logs=%d", saved, logs)
}
_, err := NewService(db).Cancel(context.Background(), m.ID, "repeat", "matched")
if err != errStateConflict {
t.Fatalf("repeat err=%v", err)
}
} else if saved.Status != tc.status || logs != 0 {
t.Fatalf("rejected transition mutated state: %+v logs=%d", saved, logs)
}
})
}
}
func TestCancelExpectedStatusMissingAndConfirmedRace(t *testing.T) {
db := testDB(t)
s := NewService(db)
if _, err := s.Cancel(context.Background(), 99999, "tester", "matched"); err != gorm.ErrRecordNotFound {
t.Fatalf("missing err=%v", err)
}
m := seedCancelMatch(t, db, "matched")
if _, err := s.Confirm(context.Background(), m.ID, "reviewer"); err != nil {
t.Fatal(err)
}
// A batch holding the earlier matched snapshot must not undo a completed confirmation.
if _, err := s.Cancel(context.Background(), m.ID, "batch", "matched"); err != errStateConflict {
t.Fatalf("stale batch err=%v", err)
}
var saved models.ReturnMatch
db.First(&saved, m.ID)
if saved.Status != "confirmed" || saved.ActiveSYBProductID == nil {
t.Fatalf("confirmation lost: %+v", saved)
}
if _, err := s.Cancel(context.Background(), m.ID, "legacy"); err != nil {
t.Fatalf("legacy confirmed cancellation: %v", err)
}
}
func TestCancelAuditFailureRollsBackTransition(t *testing.T) {
db := testDB(t)
m := seedCancelMatch(t, db, "matched")
if err := db.Callback().Create().Before("gorm:create").Register("test:reject_cancel_log", func(tx *gorm.DB) {
if tx.Statement.Schema != nil && tx.Statement.Schema.Name == "ReturnMatchLog" {
tx.AddError(errors.New("synthetic audit failure"))
}
}); err != nil {
t.Fatal(err)
}
if _, err := NewService(db).Cancel(context.Background(), m.ID, "tester", "matched"); err == nil {
t.Fatal("expected audit failure")
}
var saved models.ReturnMatch
if err := db.First(&saved, m.ID).Error; err != nil {
t.Fatal(err)
}
if saved.Status != "matched" || saved.ActiveSYBProductID == nil || saved.ActiveYeekeReturnItemID == nil || saved.CancelledAt != nil {
t.Fatalf("failed audit did not roll back: %+v", saved)
}
}
@@ -0,0 +1,103 @@
package returnmatch
import (
"context"
"errors"
"fmt"
"net/http/httptest"
"testing"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"go-admin/app/goauto/models"
"gorm.io/gorm"
)
// #349 reuses the existing single-record transaction; no production batch API.
func TestConfirmBatchExistingHTTPContract(t *testing.T) {
for _, tc := range []struct {
name, status, role string
code int
}{
{"admin", "matched", "admin", 200},
{"purchaser", "matched", "purchaser", 200},
{"after_sales", "matched", "after_sales", 200},
{"already_confirmed", "confirmed", "after_sales", 409},
{"cancelled", "cancelled", "purchaser", 409},
{"forbidden", "matched", "viewer", 403},
} {
t.Run(tc.name, func(t *testing.T) {
db := testDB(t)
m := seedCancelMatch(t, db, tc.status)
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Request = httptest.NewRequest("POST", fmt.Sprintf("/return-matches/%d/confirm", m.ID), nil)
c.Params = gin.Params{{Key: "id", Value: fmt.Sprint(m.ID)}}
c.Set("JWT_PAYLOAD", jwt.MapClaims{"rolekey": tc.role, "username": "confirm-tester"})
Handler{DB: db}.Confirm(c)
if w.Code != tc.code {
t.Fatalf("status=%d body=%s", w.Code, w.Body.String())
}
var saved models.ReturnMatch
db.First(&saved, m.ID)
var logs int64
db.Model(&models.ReturnMatchLog{}).Where("match_id = ?", m.ID).Count(&logs)
if tc.code == 200 {
if saved.Status != "confirmed" || saved.ConfirmedBy != "confirm-tester" || saved.ConfirmedAt == nil || saved.ActiveSYBProductID == nil || saved.ActiveYeekeReturnItemID == nil || logs != 1 {
t.Fatalf("confirmation audit or occupancy lost: %+v logs=%d", saved, logs)
}
} else if saved.Status != tc.status || logs != 0 {
t.Fatalf("rejected confirm changed state: %+v logs=%d", saved, logs)
}
})
}
}
func TestConfirmBatchOriginalRecordStateChanges(t *testing.T) {
db := testDB(t)
s := NewService(db)
if _, err := s.Confirm(context.Background(), 99999, "tester"); err != gorm.ErrRecordNotFound {
t.Fatalf("missing=%v", err)
}
m := seedCancelMatch(t, db, "matched")
if _, err := s.Cancel(context.Background(), m.ID, "other"); err != nil {
t.Fatal(err)
}
// Simulate a later rematch with the same product/return, but a new identity.
replacement := m
replacement.ID = 0
replacement.CancelledAt = nil
replacement.CancelledBy = ""
if err := db.Create(&replacement).Error; err != nil {
t.Fatal(err)
}
if _, err := s.Confirm(context.Background(), m.ID, "batch"); err != errStateConflict {
t.Fatalf("stale original=%v", err)
}
if _, err := s.Confirm(context.Background(), replacement.ID, "reviewer"); err != nil {
t.Fatal(err)
}
if _, err := s.Confirm(context.Background(), replacement.ID, "repeat"); err != errStateConflict {
t.Fatalf("duplicate=%v", err)
}
}
func TestConfirmBatchAuditFailureRollsBack(t *testing.T) {
db := testDB(t)
m := seedCancelMatch(t, db, "matched")
if err := db.Callback().Create().Before("gorm:create").Register("test:reject_confirm_log", func(tx *gorm.DB) {
if tx.Statement.Schema != nil && tx.Statement.Schema.Name == "ReturnMatchLog" {
tx.AddError(errors.New("synthetic audit failure"))
}
}); err != nil {
t.Fatal(err)
}
if _, err := NewService(db).Confirm(context.Background(), m.ID, "tester"); err == nil {
t.Fatal("expected audit failure")
}
var saved models.ReturnMatch
db.First(&saved, m.ID)
if saved.Status != "matched" || saved.ConfirmedAt != nil || saved.ActiveSYBProductID == nil || saved.ActiveYeekeReturnItemID == nil {
t.Fatalf("confirm did not roll back: %+v", saved)
}
}
+22 -4
View File
@@ -1,7 +1,9 @@
package returnmatch
import (
"encoding/json"
"errors"
"io"
"net/http"
"strconv"
@@ -44,14 +46,14 @@ func operatorFromContext(c *gin.Context) (string, string) {
return role, username
}
// requireCanPurchase mirrors the admin/purchaser write gate this codebase
// requireCanPurchase mirrors the admin/purchaser/after-sales write gate this codebase
// already uses for other manual-trigger actions (yeeke.Handler.TriggerSync,
// sybimport.Handler.Import): trigger match, confirm, cancel and remark are
// writes and require it; the two list/detail read endpoints do not.
func requireCanPurchase(c *gin.Context) bool {
role, _ := operatorFromContext(c)
if role != "admin" && role != "purchaser" {
c.JSON(http.StatusForbidden, gin.H{"code": "FORBIDDEN", "message": "只有管理员或采购员可以操作退货匹配"})
if role != "admin" && role != "purchaser" && role != "after_sales" {
c.JSON(http.StatusForbidden, gin.H{"code": "FORBIDDEN", "message": "只有管理员、采购员或售后可以操作退货匹配"})
return false
}
return true
@@ -140,7 +142,19 @@ func (h Handler) Confirm(c *gin.Context) {
func (h Handler) Cancel(c *gin.Context) {
h.transition(c, func(s *Service, ctx *gin.Context, id uint64, operator string) (models.ReturnMatch, error) {
return s.Cancel(ctx.Request.Context(), id, operator)
var body map[string]json.RawMessage
if err := ctx.ShouldBindJSON(&body); err != nil && !errors.Is(err, io.EOF) {
return models.ReturnMatch{}, errInvalidExpectedStatus
}
raw, supplied := body["expectedStatus"]
if !supplied {
return s.Cancel(ctx.Request.Context(), id, operator)
}
var expectedStatus string
if err := json.Unmarshal(raw, &expectedStatus); err != nil {
return models.ReturnMatch{}, errInvalidExpectedStatus
}
return s.Cancel(ctx.Request.Context(), id, operator, expectedStatus)
})
}
@@ -160,6 +174,10 @@ func (h Handler) transition(c *gin.Context, fn func(*Service, *gin.Context, uint
_, operator := operatorFromContext(c)
match, err := fn(NewService(db), c, id, operator)
if err != nil {
if errors.Is(err, errInvalidExpectedStatus) {
c.JSON(http.StatusBadRequest, gin.H{"code": "INVALID_REQUEST", "message": "expectedStatus 只允许 matched,或省略请求体"})
return
}
if errors.Is(err, gorm.ErrRecordNotFound) {
c.JSON(http.StatusNotFound, gin.H{"code": "NOT_FOUND", "message": "匹配记录不存在"})
return
@@ -0,0 +1,26 @@
package returnmatch
import (
"net/http/httptest"
"testing"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
)
func TestReturnMatchWriteRoles(t *testing.T) {
for _, role := range []string{"admin", "purchaser", "after_sales", "other", ""} {
t.Run(role, func(t *testing.T) {
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Set("JWT_PAYLOAD", jwt.MapClaims{"rolekey": role})
want := role == "admin" || role == "purchaser" || role == "after_sales"
if got := requireCanPurchase(c); got != want {
t.Fatalf("role %q allowed=%v want %v", role, got, want)
}
if !want && w.Code != 403 {
t.Fatalf("denied role status=%d", w.Code)
}
})
}
}
@@ -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")
}
}
+1 -1
View File
@@ -7,7 +7,7 @@ import (
// InitRouter mounts the #338 return-matching admin surface. List/detail are
// readable by any authenticated admin user; the write actions (batch match,
// confirm, cancel, remark) additionally require admin/purchaser via
// confirm, cancel, remark) additionally require admin/purchaser/after-sales via
// requireCanPurchase, same gate as yeeke.Handler.TriggerSync.
func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) {
handler := Handler{}
+33 -5
View File
@@ -244,6 +244,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.
@@ -298,6 +303,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
@@ -318,6 +324,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
@@ -346,6 +353,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
@@ -402,11 +420,18 @@ func (s *Service) Confirm(ctx context.Context, matchID uint64, operator string)
return match, err
}
// Cancel restores the SYB product to purchasable (by clearing
// ActiveSYBProductID) and returns the return item to the eligible pool (by
// clearing ActiveYeekeReturnItemID), from either matched or confirmed state,
// Cancel releases the active SYB product and return item pointers; actual
// purchase readiness and return-pool eligibility are then re-evaluated.
// Without a precondition it accepts either matched or confirmed state,
// per issue #338 rule: 取消匹配后再次点击「匹配退货」若配回同一对,允许.
func (s *Service) Cancel(ctx context.Context, matchID uint64, operator string) (models.ReturnMatch, error) {
var errInvalidExpectedStatus = errors.New("expectedStatus must be matched")
// Optional matched precondition is checked under the same row lock as cancellation.
// Legacy callers without it retain the confirmed-to-cancelled transition.
func (s *Service) Cancel(ctx context.Context, matchID uint64, operator string, expectedStatus ...string) (models.ReturnMatch, error) {
if len(expectedStatus) > 1 || (len(expectedStatus) == 1 && expectedStatus[0] != models.ReturnMatchStatusMatched) {
return models.ReturnMatch{}, errInvalidExpectedStatus
}
var match models.ReturnMatch
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.Clauses(clauseLockUpdate()).First(&match, matchID).Error; err != nil {
@@ -415,6 +440,9 @@ func (s *Service) Cancel(ctx context.Context, matchID uint64, operator string) (
if match.Status == models.ReturnMatchStatusCancelled {
return errStateConflict
}
if len(expectedStatus) == 1 && (match.Status != expectedStatus[0] || match.ActiveSYBProductID == nil || match.ActiveYeekeReturnItemID == nil) {
return errStateConflict
}
now := s.Now()
match.Status = models.ReturnMatchStatusCancelled
match.ActiveSYBProductID = nil
@@ -424,7 +452,7 @@ func (s *Service) Cancel(ctx context.Context, matchID uint64, operator string) (
if err := tx.Save(&match).Error; err != nil {
return err
}
return tx.Create(&models.ReturnMatchLog{MatchID: match.ID, Action: models.ReturnMatchLogActionCancelled, Operator: operator, Detail: "取消匹配,SYB 商品恢复可采购,退货商品回到可用池"}).Error
return tx.Create(&models.ReturnMatchLog{MatchID: match.ID, Action: models.ReturnMatchLogActionCancelled, Operator: operator, Detail: "取消匹配,释放 SYB 与退货商品占用,采购准备状态重新计算"}).Error
})
return match, err
}
@@ -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 {
+4
View File
@@ -30,6 +30,10 @@ type SessionStore struct{ db *gorm.DB }
func NewSessionStore(db *gorm.DB) *SessionStore { return &SessionStore{db: db} }
// DB exposes the store connection to the shared session manager; callers do
// not receive any session data through this accessor.
func (s *SessionStore) DB() *gorm.DB { return s.db }
// Session is one cached SYB login.
type Session struct {
Username string
@@ -0,0 +1,147 @@
package sybclient
import (
"context"
"errors"
"fmt"
"time"
"github.com/google/uuid"
"go-admin/app/goauto/models"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
// LoginConfig contains the non-secret connection settings and credentials
// needed to acquire a SYB session. Password is used only during LoginWithOCR.
type LoginConfig struct {
BaseURL string
Username string
Password string
OcrURL string
OcrMaxAttempts int
}
// AcquireSession reuses a valid cached session and performs at most one
// re-login per account at a time. Callers waiting for another process to
// refresh the session never request another captcha.
func AcquireSession(ctx context.Context, db *gorm.DB, cfg LoginConfig) (*Client, error) {
if cfg.Username == "" {
return nil, fmt.Errorf("顺云宝账号未配置")
}
client, err := New(cfg.BaseURL)
if err != nil {
return nil, err
}
store := NewSessionStore(db)
if cached, loadErr := store.Load(ctx, cfg.Username, time.Now()); loadErr == nil {
if importErr := client.ImportCookiesJSON(cached.CookiesJSON); importErr != nil {
return nil, importErr
}
checkErr := client.CheckSession(ctx, cached.UserID, cfg.Username)
if checkErr == nil {
return client, nil
}
if !errors.Is(checkErr, ErrSessionInvalid) {
return nil, checkErr
}
if cfg.Password == "" {
return nil, checkErr
}
} else if !errors.Is(loadErr, ErrNoSession) {
return nil, loadErr
} else if cfg.Password == "" {
return nil, loadErr
}
owner := uuid.NewString()
if claimed, err := claimAuthLease(ctx, db, cfg.Username, owner, 2*time.Minute); err != nil {
return nil, err
} else if claimed {
defer releaseAuthLease(context.Background(), db, cfg.Username, owner)
// Another worker may have completed login between our first probe and
// acquiring the lease; always re-check before requesting a captcha.
if c, ok := validCachedSession(ctx, client, store, cfg.Username); ok {
return c, nil
}
if err := refreshSession(ctx, db, client, store, cfg); err != nil {
return nil, err
}
return validSessionOrError(ctx, client, store, cfg.Username)
}
// A different worker owns the login lease. Wait for its session, bounded by
// the caller's context; do not trigger a second login.
for {
if c, ok := validCachedSession(ctx, client, store, cfg.Username); ok {
return c, nil
}
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("等待顺云宝会话刷新超时: %w", err)
}
select {
case <-time.After(500 * time.Millisecond):
case <-ctx.Done():
return nil, ctx.Err()
}
}
}
func validSessionOrError(ctx context.Context, client *Client, store *SessionStore, username string) (*Client, error) {
s, err := store.Load(ctx, username, time.Now())
if err != nil {
return nil, err
}
if err := client.ImportCookiesJSON(s.CookiesJSON); err != nil {
return nil, err
}
if err := client.CheckSession(ctx, s.UserID, username); err != nil {
return nil, err
}
return client, nil
}
func validCachedSession(ctx context.Context, client *Client, store *SessionStore, username string) (*Client, bool) {
s, err := store.Load(ctx, username, time.Now())
if err != nil {
return nil, false
}
if err := client.ImportCookiesJSON(s.CookiesJSON); err != nil {
return nil, false
}
if err := client.CheckSession(ctx, s.UserID, username); err != nil {
// Network failures are not treated as logout, but they also do not
// authorize a write; the caller will retry through the normal worker.
return nil, false
}
return client, true
}
func refreshSession(ctx context.Context, db *gorm.DB, client *Client, store *SessionStore, cfg LoginConfig) error {
ocr, err := NewOcrClient(cfg.OcrURL, 0)
if err != nil {
return fmt.Errorf("顺云宝验证码识别服务不可用: %w", err)
}
result, reason := client.LoginWithOCR(ctx, ocr, cfg.Username, cfg.Password, cfg.OcrMaxAttempts)
if result == nil {
return fmt.Errorf("顺云宝自动登录失败,需要手工输入验证码: %s", reason)
}
cookies, err := client.ExportCookiesJSON()
if err != nil {
return err
}
return store.Save(ctx, Session{Username: cfg.Username, UserID: result.User.ID, CookiesJSON: cookies, ExpiresAt: result.ExpiresAt})
}
func claimAuthLease(ctx context.Context, db *gorm.DB, username, owner string, ttl time.Duration) (bool, error) {
now := time.Now().UTC()
if err := db.WithContext(ctx).Clauses(clause.OnConflict{DoNothing: true}).Create(&models.SYBSessionAuthLease{ID: 1, Username: username, Owner: "", ExpiresAt: now.Add(-time.Second)}).Error; err != nil {
return false, err
}
r := db.WithContext(ctx).Model(&models.SYBSessionAuthLease{}).Where("id = 1 AND username = ? AND (expires_at <= ? OR owner = '')", username, now).Updates(map[string]any{"owner": owner, "expires_at": now.Add(ttl)})
return r.RowsAffected == 1, r.Error
}
func releaseAuthLease(ctx context.Context, db *gorm.DB, username, owner string) {
_ = db.WithContext(ctx).Model(&models.SYBSessionAuthLease{}).Where("id = 1 AND username = ? AND owner = ?", username, owner).Updates(map[string]any{"owner": "", "expires_at": time.Now().UTC().Add(-time.Second)})
}
@@ -0,0 +1,35 @@
package sybclient
import (
"context"
"testing"
"time"
"go-admin/app/goauto/models"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
func TestAuthLeaseAllowsOnlyOneRefreshOwner(t *testing.T) {
db, err := gorm.Open(sqlite.Open("file:syb-auth-lease?mode=memory&cache=shared"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := db.AutoMigrate(&models.SYBSessionAuthLease{}); err != nil {
t.Fatal(err)
}
ctx := context.Background()
first, err := claimAuthLease(ctx, db, "operator", "first", time.Minute)
if err != nil || !first {
t.Fatalf("first owner should claim: %v %v", first, err)
}
second, err := claimAuthLease(ctx, db, "operator", "second", time.Minute)
if err != nil || second {
t.Fatalf("second owner must wait: %v %v", second, err)
}
releaseAuthLease(ctx, db, "operator", "first")
second, err = claimAuthLease(ctx, db, "operator", "second", time.Minute)
if err != nil || !second {
t.Fatalf("lease should be reusable: %v %v", second, err)
}
}
+1 -1
View File
@@ -47,7 +47,7 @@ func (handler Handler) List(c *gin.Context) {
return
}
response, err := service.List(c.Request.Context(), ListRequest{
Page: page, PageSize: pageSize, ShopName: c.Query("shopName"), OrderCodes: []string{c.Query("orderCodes")}, ParseStatus: strings.TrimSpace(c.Query("parseStatus")), ProcessStage: strings.TrimSpace(c.Query("processStage")), PurchaseType: strings.TrimSpace(c.Query("purchaseType")),
Page: page, PageSize: pageSize, ShopName: c.Query("shopName"), OrderCodes: []string{c.Query("orderCodes")}, ParseStatus: strings.TrimSpace(c.Query("parseStatus")), ProcessStage: strings.TrimSpace(c.Query("processStage")), CreatedFrom: strings.TrimSpace(c.Query("createdFrom")), CreatedTo: strings.TrimSpace(c.Query("createdTo")), PurchaseType: strings.TrimSpace(c.Query("purchaseType")),
})
if err != nil {
writeError(c, err)
+53 -2
View File
@@ -5,6 +5,7 @@ import (
"errors"
"fmt"
"strings"
"time"
"go-admin/app/goauto/models"
"go-admin/app/goauto/purchase"
@@ -44,11 +45,14 @@ type ListRequest struct {
OrderCodes []string
ParseStatus string
ProcessStage string
CreatedFrom string
CreatedTo string
// PurchaseType is #340's list-side isolation filter: "pdd" (default when
// empty) shows only rows that still need a PDD purchase,
// "excluded" shows only pdd_purchase_excluded rows, "all" shows both. It
// combines with ProcessStage as AND; the auto-switch to 全部 mentioned in
// the issue is a front-end behaviour, not a server default.
// combines with ProcessStage and the created-time range as AND; the
// auto-switch to 全部 mentioned in the issue is a front-end behaviour,
// not a server default.
PurchaseType string
}
@@ -81,6 +85,16 @@ func (service *Service) List(ctx context.Context, request ListRequest) (ListResp
request.PageSize = 500
}
query := service.DB.WithContext(ctx).Model(&models.SYBProduct{})
createdFrom, createdTo, err := createdAtRange(request.CreatedFrom, request.CreatedTo)
if err != nil {
return ListResponse{}, err
}
if createdFrom != nil {
query = query.Where("created_at >= ?", *createdFrom)
}
if createdTo != nil {
query = query.Where("created_at < ?", *createdTo)
}
if request.ShopName = strings.TrimSpace(request.ShopName); request.ShopName != "" {
if len([]rune(request.ShopName)) > 255 {
return ListResponse{}, invalidRequest("店铺名称不能超过 255 个字符")
@@ -164,6 +178,43 @@ func (service *Service) List(ctx context.Context, request ListRequest) (ListResp
return ListResponse{Items: items, Total: total, Page: request.Page, PageSize: request.PageSize}, nil
}
// createdAtRange turns inclusive YYYY-MM-DD bounds into a half-open time
// range. The bounds are interpreted in the server's local timezone, matching
// the timestamps written by GORM for this service.
func createdAtRange(from, to string) (*time.Time, *time.Time, error) {
from = strings.TrimSpace(from)
to = strings.TrimSpace(to)
if from == "" && to == "" {
return nil, nil, nil
}
parse := func(value, label string) (*time.Time, error) {
if value == "" {
return nil, nil
}
parsed, err := time.ParseInLocation("2006-01-02", value, time.Local)
if err != nil {
return nil, invalidRequest(label + " 必须是 YYYY-MM-DD")
}
return &parsed, nil
}
start, err := parse(from, "createdFrom")
if err != nil {
return nil, nil, err
}
endDay, err := parse(to, "createdTo")
if err != nil {
return nil, nil, err
}
if start != nil && endDay != nil && start.After(*endDay) {
return nil, nil, invalidRequest("createdFrom 不能晚于 createdTo")
}
if endDay != nil {
end := endDay.AddDate(0, 0, 1)
endDay = &end
}
return start, endDay, nil
}
func normalizeOrderCodes(raw []string) ([]string, error) {
seen := make(map[string]bool, len(raw))
result := make([]string, 0, len(raw))
@@ -7,6 +7,7 @@ import (
"fmt"
"strings"
"testing"
"time"
"go-admin/app/goauto/models"
"go-admin/app/goauto/sybimport"
@@ -123,6 +124,46 @@ func TestServiceListCapsPageSizeAt500(t *testing.T) {
}
}
func TestServiceListFiltersByCreatedDateInclusive(t *testing.T) {
db := openTestDB(t)
order := realOrder()
for i, day := range []string{"2026-09-01", "2026-09-02", "2026-09-03"} {
rowOrder := order
rowOrder.Code = fmt.Sprintf("CREATED-%d", i)
rowOrder.StockID += uint64(i)
detail := realDetailA()
detail.ID += uint64(i)
result, err := sybimport.ApplyDetail(context.Background(), db, rowOrder, detail)
if err != nil {
t.Fatal(err)
}
created, err := time.ParseInLocation("2006-01-02", day, time.Local)
if err != nil {
t.Fatal(err)
}
if err := db.Model(&models.SYBProduct{}).Where("id = ?", result.SYBProduct.ID).Update("created_at", created).Error; err != nil {
t.Fatal(err)
}
}
service := sybimport.NewService(db)
between, err := service.List(context.Background(), sybimport.ListRequest{CreatedFrom: "2026-09-01", CreatedTo: "2026-09-02"})
if err != nil || between.Total != 2 {
t.Fatalf("inclusive created date range should return two rows, total=%d err=%v", between.Total, err)
}
fromOnly, err := service.List(context.Background(), sybimport.ListRequest{CreatedFrom: "2026-09-03"})
if err != nil || fromOnly.Total != 1 {
t.Fatalf("created-from filter should return one row, total=%d err=%v", fromOnly.Total, err)
}
toOnly, err := service.List(context.Background(), sybimport.ListRequest{CreatedTo: "2026-09-01"})
if err != nil || toOnly.Total != 1 {
t.Fatalf("created-to filter should return one row, total=%d err=%v", toOnly.Total, err)
}
if _, err := service.List(context.Background(), sybimport.ListRequest{CreatedFrom: "2026-09-04", CreatedTo: "2026-09-01"}); serviceErrCode(t, err) != sybimport.CodeInvalidRequest {
t.Fatalf("reversed created date range should be rejected: %v", err)
}
}
func TestServiceListRejectsInvalidParseStatus(t *testing.T) {
db := openTestDB(t)
service := sybimport.NewService(db)
+4 -48
View File
@@ -491,54 +491,10 @@ func splitDateRange(dateFrom, dateTo string) ([]string, error) {
// blip as a logout would trigger needless logins and could throw away a
// perfectly good session.
func Connect(ctx context.Context, store *sybclient.SessionStore, cfg ConnectConfig) (*sybclient.Client, error) {
if cfg.Username == "" || cfg.Password == "" {
return nil, errors.New("顺云宝账号或密码未配置,请设置 GOAUTO_SYB_USERNAME 和 GOAUTO_SYB_PASSWORD")
}
client, err := sybclient.New(cfg.BaseURL)
if err != nil {
return nil, err
}
cached, err := store.Load(ctx, cfg.Username, time.Now())
switch {
case err == nil:
if importErr := client.ImportCookiesJSON(cached.CookiesJSON); importErr == nil {
if checkErr := client.CheckSession(ctx, cached.UserID, cfg.Username); checkErr == nil {
return client, nil
} else if errors.Is(checkErr, sybclient.ErrSessionInvalid) {
if delErr := store.Delete(ctx, cfg.Username); delErr != nil {
return nil, delErr
}
}
}
case errors.Is(err, sybclient.ErrNoSession):
// Nothing cached; fall through to a fresh login.
default:
return nil, err
}
if cfg.OcrURL == "" {
return nil, errors.New("顺云宝会话已失效,且未配置验证码识别服务;请配置 extend.syb.ocrurl 或改用手工登录")
}
ocr, err := sybclient.NewOcrClient(cfg.OcrURL, 0)
if err != nil {
return nil, err
}
result, reason := client.LoginWithOCR(ctx, ocr, cfg.Username, cfg.Password, cfg.OcrMaxAttempts)
if result == nil {
return nil, fmt.Errorf("顺云宝自动登录失败,需要手工输入验证码: %s", reason)
}
jar, err := client.ExportCookiesJSON()
if err != nil {
return nil, err
}
if err := store.Save(ctx, sybclient.Session{
Username: cfg.Username, UserID: result.User.ID,
CookiesJSON: jar, ExpiresAt: result.ExpiresAt,
}); err != nil {
return nil, err
}
return client, nil
return sybclient.AcquireSession(ctx, store.DB(), sybclient.LoginConfig{
BaseURL: cfg.BaseURL, Username: cfg.Username, Password: cfg.Password,
OcrURL: cfg.OcrURL, OcrMaxAttempts: cfg.OcrMaxAttempts,
})
}
// ConnectConfig carries the login settings. The password is passed through and
+206 -5
View File
@@ -10,6 +10,7 @@ import (
"sort"
"strings"
"time"
"unicode"
"go-admin/app/goauto/models"
"go-admin/app/goauto/sybclient"
@@ -193,13 +194,20 @@ func planRecord(ctx context.Context, reader MatchReader, record models.SYBInnerC
if len(matches) == 1 && matches[0].ProductQty == count {
chosen = []sybclient.DetailItem{matches[0]}
} else if len(matches) == count && count > 1 {
sort.Slice(matches, func(i, j int) bool { return matches[i].ID < matches[j].ID })
for _, item := range matches {
if item.ProductQty != 1 {
return nil, models.SYBInnerCodeSkipped, "相同规格候选数量不明确,不能自动分配", nil
}
}
chosen = matches
codes := make([]string, count)
for i, it := range record.Items {
codes[i] = it.Code
}
assigned, reason := assignExistingBoundItems(record.Stall, codes, matches)
if reason != "" {
return nil, models.SYBInnerCodeSkipped, reason, nil
}
chosen = assigned
} else if len(matches) > 1 {
return nil, models.SYBInnerCodeSkipped, "同一订单存在多条相同规格候选商品,不能自动选择", nil
} else {
@@ -246,6 +254,83 @@ func planRecord(ctx context.Context, reader MatchReader, record models.SYBInnerC
return plan, models.SYBInnerCodeReady, "唯一匹配,等待确认回写", nil
}
// assignExistingBoundItems 把 N 个待写入入库码按顺序分配给 N 个数量为 1 的候选商品明细。
// 修复 #289:候选明细的匹配顺序(按 ID 排序)未必与目标码顺序一致,若单纯按下标
// 对应,会把已经正确绑定某个目标码的明细错误地重新分配给另一个码。这里先按“候选
// 明细已有的入库码值”精确匹配对应的目标码,保留既有正确绑定不动;再把剩余尚未
// 写入任何码的空白明细(按 ID 排序)依次填充给还没有候选的目标码位置。
// 移植自 cmautobuy `planExistingMatchedInnerCodeItems` 的一致性护栏(代码评审补充):
// 候选来自 NormalizeSpecKey 归一化匹配,原始 ProductSpec/sku/variationSku 可能在
// 归一化后相同但原始值不同,必须逐一比对最低 ID 候选,避免跨真正不同商品自动分配;
// 无 SKU 回退路径可能返回从未做过档口校验的候选,这里逐一重新校验;同时拒绝无效
// 或重复的商品明细 ID。
func assignExistingBoundItems(stall string, codes []string, matches []sybclient.DetailItem) ([]sybclient.DetailItem, string) {
sorted := append([]sybclient.DetailItem(nil), matches...)
sort.Slice(sorted, func(i, j int) bool { return sorted[i].ID < sorted[j].ID })
seenIDs := make(map[int64]bool, len(sorted))
for _, item := range sorted {
if item.ID <= 0 || seenIDs[item.ID] {
return nil, "重复候选包含无效或重复的商品明细 ID,不能自动逐件分配"
}
seenIDs[item.ID] = true
}
first := sorted[0]
wantSpec := first.ProductSpec
wantSKU := rawText(first.Raw["sku"])
wantVariationSKU := rawText(first.Raw["variationSku"])
for _, item := range sorted {
if item.ProductSpec != wantSpec || rawText(item.Raw["sku"]) != wantSKU || rawText(item.Raw["variationSku"]) != wantVariationSKU {
return nil, "重复候选的规格或 SKU 身份不一致,不能自动逐件分配"
}
}
if strings.TrimSpace(stall) != "" {
for _, item := range sorted {
if !stallMatches(stall, item) {
return nil, "重复候选的档口及货号不一致,不能自动逐件分配"
}
}
}
assigned := make([]sybclient.DetailItem, len(codes))
taken := make([]bool, len(codes))
codeIndex := make(map[string]int, len(codes))
for i, code := range codes {
codeIndex[code] = i
}
blanks := make([]int, 0, len(sorted))
for si, item := range sorted {
remote := rawText(item.Raw["innerExpCode"])
if remote == "" {
blanks = append(blanks, si)
continue
}
idx, ok := codeIndex[remote]
if !ok {
return nil, "候选商品明细存在非目标入库码,不能自动逐件分配"
}
if taken[idx] {
return nil, "同一入库码在候选商品中出现多次,不能自动逐件分配"
}
assigned[idx] = item
taken[idx] = true
}
bi := 0
for i := range assigned {
if taken[i] {
continue
}
if bi >= len(blanks) {
return nil, "现成空白明细不足,不能完成逐件分配"
}
assigned[i] = sorted[blanks[bi]]
taken[i] = true
bi++
}
if bi != len(blanks) {
return nil, "现成空白明细多于待写入入库码,不能自动逐件分配"
}
return assigned, ""
}
func matchSpec(spec string, items []sybclient.DetailItem) []sybclient.DetailItem {
result := []sybclient.DetailItem{}
for _, item := range items {
@@ -305,15 +390,131 @@ func strictStall(stall string, items []sybclient.DetailItem) []sybclient.DetailI
return nil
}
result := []sybclient.DetailItem{}
name, article, has := strings.Cut(stall, "#")
for _, item := range items {
blob := rawText(item.Raw["sku"]) + " " + rawText(item.Raw["variationSku"]) + " " + item.ProductSpec
if strings.Contains(blob, stall) || (has && strings.Contains(blob, strings.TrimSpace(name)) && strings.Contains(blob, strings.TrimSpace(article))) {
if stallMatches(stall, item) {
result = append(result, item)
}
}
return result
}
// stallMatches 移植自 cmautobuy `innerCodeStallMatches`(#259/#273 修复):
// 档口名与货号以最后一个 `#` 切分;货号只与字母数字 token 比较;纯数字货号要求
// 候选中同时包含档口名才允许前导零等价(如 "067"≡"67");非数字货号要求精确
// token 匹配;ProductSpec 只在以货号开头时才算命中;货号为空时回退为档口名包含判断。
func stallMatches(stall string, item sybclient.DetailItem) bool {
sku := rawText(item.Raw["sku"])
variation := rawText(item.Raw["variationSku"])
blob := sku + " " + variation + " " + item.ProductSpec
if strings.Contains(blob, stall) {
return true
}
name, article, hasArticle := splitStall(stall)
if !hasArticle {
return false
}
if article == "" {
return name != "" && (strings.Contains(sku, name) || strings.Contains(variation, name))
}
if isNumericArticle(article) {
nameMatches := name != "" && (strings.Contains(sku, name) || strings.Contains(variation, name))
if !nameMatches {
return false
}
return textHasNumericArticle(sku, article) ||
textHasNumericArticle(variation, article) ||
productSpecStartsWithArticle(item.ProductSpec, article, true)
}
return textHasExactArticle(sku, article) ||
textHasExactArticle(variation, article) ||
productSpecStartsWithArticle(item.ProductSpec, article, false)
}
// splitStall 从档口名称#货号取最后一个 #,避免档口名称本身含 # 时截错。
func splitStall(stall string) (name, article string, ok bool) {
stall = strings.TrimSpace(stall)
separator := strings.LastIndex(stall, "#")
if separator < 0 {
return stall, "", false
}
return strings.TrimSpace(stall[:separator]), strings.TrimSpace(stall[separator+1:]), true
}
func isNumericArticle(article string) bool {
if article == "" {
return false
}
for _, char := range article {
if !unicode.IsDigit(char) {
return false
}
}
return true
}
func textHasNumericArticle(text, article string) bool {
target := normalizeNumericArticle(article)
for _, token := range articleTokens(text) {
if isNumericArticle(token) && normalizeNumericArticle(token) == target {
return true
}
}
return false
}
func normalizeNumericArticle(article string) string {
normalized := strings.TrimLeft(article, "0")
if normalized == "" {
return "0"
}
return normalized
}
func textHasExactArticle(text, article string) bool {
for _, token := range articleTokens(text) {
if token == article {
return true
}
}
return false
}
func productSpecStartsWithArticle(productSpec, article string, numeric bool) bool {
productSpec = strings.TrimSpace(productSpec)
separator := strings.IndexAny(productSpec, " ,,")
if separator <= 0 {
return false
}
prefix := strings.TrimSpace(productSpec[:separator])
if numeric {
return isNumericArticle(prefix) && normalizeNumericArticle(prefix) == normalizeNumericArticle(article)
}
return prefix == article
}
// articleTokens 只把连续字母或数字视为货号候选,标点、【】、#、横线、空格自然成为
// 边界:能识别 "067【档口】",又不会把 "PDD256437" 中间的数字误认为独立货号。
func articleTokens(text string) []string {
tokens := make([]string, 0)
start := -1
runes := []rune(text)
for index, char := range runes {
if unicode.IsLetter(char) || unicode.IsDigit(char) {
if start < 0 {
start = index
}
continue
}
if start >= 0 {
tokens = append(tokens, string(runes[start:index]))
start = -1
}
}
if start >= 0 {
tokens = append(tokens, string(runes[start:]))
}
return tokens
}
func rawText(value any) string {
switch v := value.(type) {
case string:
@@ -147,6 +147,209 @@ func TestRunMatchJobReservesDifferentDetailsForRecordsInOneBatch(t *testing.T) {
}
}
// --- #259/#273/#289 stall-matching regression tests (ported from cmautobuy) ---
func TestStallMatchesIgnoresWeightLikeNumberAsArticle(t *testing.T) {
// #273: "50公斤" must not be treated as if the article were the bare number 50.
item := detail(1, "50公斤,黑色", 1, "SKU-X", "", "")
if stallMatches("档口甲#50", item) {
t.Fatalf("weight-like text must not match numeric article 50")
}
}
func TestStallMatchesRejectsSubstringInsideLongCode(t *testing.T) {
// #273: a long code with internal digits (PDD256437) must not spuriously
// match a short numeric article (256) via substring containment.
item := detail(1, "黑色,L", 1, "PDD256437", "档口甲", "")
if stallMatches("档口甲#256", item) {
t.Fatalf("long code must not match numeric article 256 via substring")
}
}
func TestStallMatchesLeadingZeroEquivalenceRequiresStallName(t *testing.T) {
// #259/#273: "067" and "67" are equivalent articles only when the stall
// name also matches; a different stall name must not match.
sameStall := detail(1, "黑色,L", 1, "档口甲-067", "", "")
if !stallMatches("档口甲#67", sameStall) {
t.Fatalf("067 should be treated as equivalent to 67 when stall name matches")
}
differentStall := detail(2, "黑色,L", 1, "档口乙-067", "", "")
if stallMatches("档口甲#67", differentStall) {
t.Fatalf("067 must not match 67 when the stall name differs")
}
}
func TestStallMatchesNonNumericArticleRequiresExactToken(t *testing.T) {
item := detail(1, "黑色,L", 1, "ABC12", "", "")
if stallMatches("档口甲#AB", item) {
t.Fatalf("non-numeric article must require an exact token match, not substring")
}
exact := detail(2, "黑色,L", 1, "AB", "", "")
if !stallMatches("档口甲#AB", exact) {
t.Fatalf("exact non-numeric token should match")
}
}
func TestStallMatchesHandlesHashInsideStallName(t *testing.T) {
// Splits on the LAST '#' so a stall name that itself contains '#' still
// yields the correct article.
item := detail(1, "黑色,L", 1, "档口#甲-67", "", "")
if !stallMatches("档口#甲#67", item) {
t.Fatalf("stall name containing '#' should still resolve article via last '#'")
}
}
func TestStallMatchesProductSpecPrefixMatchesArticle(t *testing.T) {
// Non-numeric article: ProductSpec prefix match does not additionally
// require the stall name to appear in sku/variationSku.
item := detail(1, "AB 黑色,L", 1, "", "", "")
if !stallMatches("档口甲#AB", item) {
t.Fatalf("ProductSpec starting with the article should match")
}
notPrefix := detail(2, "黑色,ABL", 1, "", "", "")
if stallMatches("档口甲#AB", notPrefix) {
t.Fatalf("article appearing mid-spec (not as prefix) must not match")
}
}
func TestStallMatchesEmptyArticleFallsBackToStallName(t *testing.T) {
item := detail(1, "黑色,L", 1, "档口甲专柜", "", "")
if !stallMatches("档口甲#", item) {
t.Fatalf("empty article should fall back to stall-name containment")
}
}
// --- #289: existing-binding-preserving multi-piece assignment ---
func TestAssignExistingBoundItemsPreservesOutOfOrderBindings(t *testing.T) {
// Two single-piece candidates already carry codes, but the previously
// bound code (IC-2) sits on the LOWER-ID detail while the target order
// expects it second. A naive ID-order/index assignment would strip the
// existing correct binding from detail 20 and try to overwrite it.
d20 := detail(20, "黑色,L", 1, "SKU-1", "A#1", "IC-2")
d21 := detail(21, "黑色,L", 1, "SKU-1", "A#1", "")
assigned, reason := assignExistingBoundItems("A#1", []string{"IC-1", "IC-2"}, []sybclient.DetailItem{d20, d21})
if reason != "" {
t.Fatalf("unexpected reason: %s", reason)
}
if assigned[1].ID != 20 {
t.Fatalf("expected detail 20 (already bound to IC-2) preserved at index 1, got %+v", assigned[1])
}
if assigned[0].ID != 21 {
t.Fatalf("expected the blank detail 21 filled in at index 0, got %+v", assigned[0])
}
}
func TestAssignExistingBoundItemsRejectsForeignCode(t *testing.T) {
d20 := detail(20, "黑色,L", 1, "SKU-1", "A#1", "IC-OTHER")
d21 := detail(21, "黑色,L", 1, "SKU-1", "A#1", "")
_, reason := assignExistingBoundItems("A#1", []string{"IC-1", "IC-2"}, []sybclient.DetailItem{d20, d21})
if reason == "" {
t.Fatalf("expected rejection for detail already holding a non-target code")
}
}
func TestAssignExistingBoundItemsRejectsInconsistentIdentity(t *testing.T) {
// Candidates can share a NormalizeSpecKey-normalized spec while their raw
// ProductSpec/sku/variationSku differ; auto-assignment across genuinely
// different items must be rejected.
d20 := detail(20, "黑色, L", 1, "SKU-1", "A#1", "")
d21 := detail(21, "黑色,L", 1, "SKU-2", "A#1", "")
_, reason := assignExistingBoundItems("A#1", []string{"IC-1", "IC-2"}, []sybclient.DetailItem{d20, d21})
if reason != "重复候选的规格或 SKU 身份不一致,不能自动逐件分配" {
t.Fatalf("expected identity-mismatch rejection, got %q", reason)
}
}
func TestAssignExistingBoundItemsRejectsCandidateFailingStallCheck(t *testing.T) {
// The no-SKU fallback path in matchEvidence can hand back candidates that
// were never checked against the stall at all. Both candidates share an
// identical identity (so the identity guard passes) but neither one's
// sku/spec actually satisfies the record's stall/article requirement.
d20 := detail(20, "黑色,L", 1, "SKU-1", "ZZZ", "")
d21 := detail(21, "黑色,L", 1, "SKU-1", "ZZZ", "")
_, reason := assignExistingBoundItems("甲档口#88", []string{"IC-1", "IC-2"}, []sybclient.DetailItem{d20, d21})
if reason != "重复候选的档口及货号不一致,不能自动逐件分配" {
t.Fatalf("expected stall-mismatch rejection, got %q", reason)
}
}
func TestAssignExistingBoundItemsRejectsInvalidOrDuplicateDetailID(t *testing.T) {
invalidID := detail(0, "黑色,L", 1, "SKU-1", "A#1", "")
valid := detail(21, "黑色,L", 1, "SKU-1", "A#1", "")
if _, reason := assignExistingBoundItems("A#1", []string{"IC-1", "IC-2"}, []sybclient.DetailItem{invalidID, valid}); reason == "" {
t.Fatalf("expected rejection for non-positive detail ID")
}
dup1 := detail(20, "黑色,L", 1, "SKU-1", "A#1", "")
dup2 := detail(20, "黑色,L", 1, "SKU-1", "A#1", "")
if _, reason := assignExistingBoundItems("A#1", []string{"IC-1", "IC-2"}, []sybclient.DetailItem{dup1, dup2}); reason == "" {
t.Fatalf("expected rejection for duplicate detail IDs")
}
}
func TestAssignExistingBoundItemsHappyPathStillAssignsWithGuards(t *testing.T) {
// Both candidates share identical raw spec/sku/variationSku and both
// individually satisfy the stall check; the guards must not block the
// legitimate happy path.
d20 := detail(20, "黑色,L", 1, "SKU-1", "A#1", "IC-2")
d21 := detail(21, "黑色,L", 1, "SKU-1", "A#1", "")
assigned, reason := assignExistingBoundItems("A#1", []string{"IC-1", "IC-2"}, []sybclient.DetailItem{d20, d21})
if reason != "" {
t.Fatalf("unexpected reason: %s", reason)
}
if assigned[0].ID != 21 || assigned[1].ID != 20 {
t.Fatalf("assigned=%+v", assigned)
}
}
func TestRunMatchJobPreservesExistingBindingWhenCandidateOrderDiffers(t *testing.T) {
// End-to-end regression for #289: N=2 single-piece candidates already
// carrying one previously bound code, with the bound detail's ID not
// matching sequential/ID order relative to the target codes. The plan
// must leave the already-correct binding untouched and only place the
// missing code onto the still-blank detail.
db := testDB(t)
// SourceSKURaw is intentionally left blank: with the same source SKU on
// both candidate rows the source-SKU path would reject as "duplicate"
// before ever reaching the stall-based multi-item assignment being
// regression-tested here.
record := models.SYBInnerCodeRecord{BusinessDate: "2026-08-28", OrderNumber: "ORDER-1", Stall: "A#1", SpecKey: "黑色,L", SpecRaw: "黑色,L", Status: models.SYBInnerCodePending, CreatedBy: 1, ImportRequestID: uuid.NewString(), Items: []models.SYBInnerCodeItem{{BusinessDate: "2026-08-28", Code: "IC-1", Ordinal: 1, SourceRow: 2}, {BusinessDate: "2026-08-28", Code: "IC-2", Ordinal: 2, SourceRow: 3}}}
records := []models.SYBInnerCodeRecord{record}
jobID := createMatchJob(t, db, records)
recordID := records[0].ID
// Detail 20 (lower ID) already carries IC-2 (bound out of sequence);
// detail 21 (higher ID) is still blank and should receive IC-1.
stock := sybclient.StockDetail{ID: 10, Details: []sybclient.DetailItem{
detail(20, "黑色,L", 1, "SKU-1", "A#1", "IC-2"),
detail(21, "黑色,L", 1, "SKU-1", "A#1", ""),
}}
reader := &fakeMatchReader{rows: map[string][]sybclient.StockRow{"ORDER-1": {{ID: 10, Code: "ORDER-1"}}}, stocks: map[int64]sybclient.StockDetail{10: stock}}
if err := RunMatchJob(context.Background(), db, reader, jobID); err != nil {
t.Fatal(err)
}
var plan models.SYBInnerCodePlan
if err := db.First(&plan, "record_id = ?", recordID).Error; err != nil {
t.Fatal(err)
}
var items []plannedRemoteItem
if err := json.Unmarshal([]byte(plan.RemoteItemsJSON), &items); err != nil {
t.Fatal(err)
}
if len(items) != 2 {
t.Fatalf("items=%+v", items)
}
byCode := map[string]plannedRemoteItem{}
for _, item := range items {
byCode[item.Code] = item
}
if byCode["IC-2"].DetailID != 20 {
t.Fatalf("existing binding for IC-2 must stay on detail 20, got %+v", byCode["IC-2"])
}
if byCode["IC-1"].DetailID != 21 {
t.Fatalf("missing IC-1 should be assigned to the blank detail 21, got %+v", byCode["IC-1"])
}
}
type captureStarter struct {
jobID string
committed bool
+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)
}
}
@@ -0,0 +1,138 @@
package version_local
import (
"errors"
"fmt"
"os"
"runtime"
"strings"
adminmodels "go-admin/app/admin/models"
"go-admin/app/goauto/access"
"go-admin/cmd/migrate/migration"
migrationmodels "go-admin/cmd/migrate/migration/models"
common "go-admin/common/models"
"gorm.io/gorm"
)
const afterSalesInitialPasswordEnv = "GOAUTO_AFTER_SALES_INITIAL_PASSWORD"
var afterSalesUsers = []string{"zengyt", "huangyj", "zhuyt", "wangxy"}
func init() {
_, fileName, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateAfterSalesRole)
}
// migrateAfterSalesRole creates the internal 售后 role and accounts. It
// copies the current purchaser menu/API grants, so the two roles stay aligned
// for the current product surface (including yeeke 退货同步 and 退货商品).
// The initial password is deliberately supplied only at migration time via an
// environment variable; it is never stored in source, logs, or issue text.
func migrateAfterSalesRole(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
role, err := ensureAfterSalesRole(tx)
if err != nil {
return err
}
if err := clonePurchaserPermissions(tx, role.RoleId); err != nil {
return err
}
if err := ensureAfterSalesUsers(tx, role.RoleId); err != nil {
return err
}
return tx.Create(&common.Migration{Version: version}).Error
})
}
func ensureAfterSalesRole(db *gorm.DB) (migrationmodels.SysRole, error) {
var purchaser migrationmodels.SysRole
if err := db.Where("role_key = ?", access.RolePurchaser).First(&purchaser).Error; err != nil {
return migrationmodels.SysRole{}, fmt.Errorf("find purchaser role: %w", err)
}
role := migrationmodels.SysRole{}
if err := db.Where("role_key = ?", access.RoleAfterSales).
Assign(migrationmodels.SysRole{
RoleName: "售后", Status: "2", RoleSort: purchaser.RoleSort + 1,
Admin: false, DataScope: purchaser.DataScope,
Remark: "GoAuto 售后业务角色(系统维护)",
}).FirstOrCreate(&role, migrationmodels.SysRole{RoleKey: access.RoleAfterSales}).Error; err != nil {
return migrationmodels.SysRole{}, err
}
return role, nil
}
func clonePurchaserPermissions(db *gorm.DB, roleID int) error {
var purchaser migrationmodels.SysRole
if err := db.Where("role_key = ?", access.RolePurchaser).First(&purchaser).Error; err != nil {
return err
}
if err := db.Exec(`
INSERT INTO sys_role_menu (role_id, menu_id)
SELECT ?, source.menu_id
FROM sys_role_menu AS source
WHERE source.role_id = ?
AND NOT EXISTS (
SELECT 1 FROM sys_role_menu AS target
WHERE target.role_id = ? AND target.menu_id = source.menu_id
)`, roleID, purchaser.RoleId, roleID).Error; err != nil {
return fmt.Errorf("clone purchaser menus: %w", err)
}
if err := db.Exec(`
INSERT INTO casbin_rule (ptype, v0, v1, v2, v3, v4, v5)
SELECT source.ptype, ?, source.v1, source.v2, source.v3, source.v4, source.v5
FROM casbin_rule AS source
WHERE source.ptype = 'p' AND source.v0 = ?
AND NOT EXISTS (
SELECT 1 FROM casbin_rule AS target
WHERE target.ptype = source.ptype AND target.v0 = ?
AND target.v1 = source.v1 AND target.v2 = source.v2
AND COALESCE(target.v3, '') = COALESCE(source.v3, '')
AND COALESCE(target.v4, '') = COALESCE(source.v4, '')
AND COALESCE(target.v5, '') = COALESCE(source.v5, '')
)`, access.RoleAfterSales, access.RolePurchaser, access.RoleAfterSales).Error; err != nil {
return fmt.Errorf("clone purchaser API policies: %w", err)
}
return nil
}
func ensureAfterSalesUsers(db *gorm.DB, roleID int) error {
missing := make([]string, 0, len(afterSalesUsers))
for _, username := range afterSalesUsers {
var existing adminmodels.SysUser
err := db.Unscoped().Where("username = ?", username).First(&existing).Error
switch {
case errors.Is(err, gorm.ErrRecordNotFound):
missing = append(missing, username)
case err != nil:
return fmt.Errorf("find after-sales user %q: %w", username, err)
case existing.RoleId != roleID:
return fmt.Errorf("after-sales user %q already exists with another role", username)
}
}
if len(missing) == 0 {
return nil
}
password := strings.TrimSpace(os.Getenv(afterSalesInitialPasswordEnv))
if password == "" {
return fmt.Errorf("%s must be set when creating after-sales users", afterSalesInitialPasswordEnv)
}
// Validate before creating any account; the value itself is never logged or
// written to a repository artifact.
if len(password) < 1 {
return errors.New("after-sales initial password is empty")
}
for _, username := range missing {
user := adminmodels.SysUser{
Username: username, Password: password, NickName: username,
RoleId: roleID, Status: "2", Remark: "GoAuto 售后账号",
}
if err := db.Create(&user).Error; err != nil {
return fmt.Errorf("create after-sales user %q: %w", username, err)
}
}
return nil
}
@@ -0,0 +1,87 @@
package version_local
import (
"os"
"testing"
adminmodels "go-admin/app/admin/models"
"go-admin/app/goauto/access"
migrationmodels "go-admin/cmd/migrate/migration/models"
"golang.org/x/crypto/bcrypt"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
func TestEnsureAfterSalesRoleAndUsersIsIdempotent(t *testing.T) {
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&migrationmodels.SysRole{}, &migrationmodels.SysMenu{}, &adminmodels.SysUser{}); err != nil {
t.Fatal(err)
}
if err = db.Exec(`CREATE TABLE IF NOT EXISTS sys_role_menu (role_id integer NOT NULL, menu_id integer NOT NULL, PRIMARY KEY (role_id, menu_id))`).Error; err != nil {
t.Fatal(err)
}
if err = db.Exec(`CREATE TABLE IF NOT EXISTS casbin_rule (id integer PRIMARY KEY AUTOINCREMENT, ptype varchar(100), v0 varchar(100), v1 varchar(100), v2 varchar(100), v3 varchar(100), v4 varchar(100), v5 varchar(100))`).Error; err != nil {
t.Fatal(err)
}
purchaser := migrationmodels.SysRole{RoleName: "采购员", RoleKey: access.RolePurchaser, Status: "2", RoleSort: 20, DataScope: "1"}
if err = db.Create(&purchaser).Error; err != nil {
t.Fatal(err)
}
if err = db.Exec(`INSERT INTO sys_menu (menu_name, title, menu_type, parent_id) VALUES ('GoAutoYeekeReturns', 'yeeke 退货包裹', 'C', 0)`).Error; err != nil {
t.Fatal(err)
}
if err = db.Exec(`INSERT INTO sys_role_menu (role_id, menu_id) SELECT ?, menu_id FROM sys_menu`, purchaser.RoleId).Error; err != nil {
t.Fatal(err)
}
if err = db.Exec(`INSERT INTO casbin_rule (ptype, v0, v1, v2) VALUES ('p', ?, '/api/admin/v1/yeeke-returns', 'GET')`, access.RolePurchaser).Error; err != nil {
t.Fatal(err)
}
oldPassword := os.Getenv(afterSalesInitialPasswordEnv)
defer os.Setenv(afterSalesInitialPasswordEnv, oldPassword)
if err = os.Setenv(afterSalesInitialPasswordEnv, "test-only-password"); err != nil {
t.Fatal(err)
}
role, err := ensureAfterSalesRole(db)
if err != nil {
t.Fatal(err)
}
if err = clonePurchaserPermissions(db, role.RoleId); err != nil {
t.Fatal(err)
}
if err = ensureAfterSalesUsers(db, role.RoleId); err != nil {
t.Fatal(err)
}
if err = clonePurchaserPermissions(db, role.RoleId); err != nil {
t.Fatal(err)
}
if err = ensureAfterSalesUsers(db, role.RoleId); err != nil {
t.Fatal(err)
}
var count int64
if err = db.Model(&adminmodels.SysUser{}).Where("role_id = ?", role.RoleId).Count(&count).Error; err != nil {
t.Fatal(err)
}
if count != int64(len(afterSalesUsers)) {
t.Fatalf("got %d after-sales users, want %d", count, len(afterSalesUsers))
}
var user adminmodels.SysUser
if err = db.Where("username = ?", afterSalesUsers[0]).First(&user).Error; err != nil {
t.Fatal(err)
}
if bcrypt.CompareHashAndPassword([]byte(user.Password), []byte("test-only-password")) != nil {
t.Fatal("initial password was not stored as a bcrypt hash")
}
if err = db.Table("casbin_rule").Where("ptype = 'p' AND v0 = ?", access.RoleAfterSales).Count(&count).Error; err != nil {
t.Fatal(err)
}
if count != 1 {
t.Fatalf("got %d after-sales policies, want 1", count)
}
}
@@ -0,0 +1,76 @@
package version_local
import (
"fmt"
"runtime"
"go-admin/app/goauto/access"
"go-admin/cmd/migrate/migration"
common "go-admin/common/models"
"gorm.io/gorm"
)
// #341 follow-up: the after-sales role was created before #338 added the
// return-match API surface. Add only the reviewed return-match grants so the
// existing role can use the new SYB action without recreating or changing its
// users.
func init() {
_, fileName, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateAfterSalesReturnMatch)
}
type afterSalesReturnMatchAPI struct {
ID int `gorm:"column:id;primaryKey;autoIncrement"`
Title string `gorm:"column:title;size:128"`
Path string `gorm:"column:path;size:128"`
Type string `gorm:"column:type;size:16"`
Action string `gorm:"column:action;size:16"`
}
func (afterSalesReturnMatchAPI) TableName() string { return "sys_api" }
type afterSalesReturnMatchPolicy struct {
ID uint `gorm:"column:id;primaryKey;autoIncrement"`
Ptype string `gorm:"column:ptype;size:100"`
V0 string `gorm:"column:v0;size:100"`
V1 string `gorm:"column:v1;size:100"`
V2 string `gorm:"column:v2;size:100"`
V3 string `gorm:"column:v3;size:100"`
V4 string `gorm:"column:v4;size:100"`
V5 string `gorm:"column:v5;size:100"`
}
func (afterSalesReturnMatchPolicy) TableName() string { return "casbin_rule" }
func migrateAfterSalesReturnMatch(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
var role struct {
RoleID int `gorm:"column:role_id"`
}
if err := tx.Table("sys_role").Select("role_id").Where("role_key = ?", access.RoleAfterSales).First(&role).Error; err != nil {
return fmt.Errorf("find after-sales role: %w", err)
}
for _, permission := range access.AdminAPIs {
if len(permission.Path) < len("/api/admin/v1/return-matches") || permission.Path[:len("/api/admin/v1/return-matches")] != "/api/admin/v1/return-matches" {
continue
}
api := afterSalesReturnMatchAPI{Path: permission.Path, Action: permission.Method}
if err := tx.Where("path = ? AND action = ?", permission.Path, permission.Method).
Assign(afterSalesReturnMatchAPI{Title: permission.Title, Path: permission.Path, Action: permission.Method, Type: "BUS"}).
FirstOrCreate(&api).Error; err != nil {
return fmt.Errorf("ensure return-match API: %w", err)
}
policy := afterSalesReturnMatchPolicy{Ptype: "p", V0: access.RoleAfterSales, V1: permission.Path, V2: permission.Method}
var count int64
if err := tx.Model(&afterSalesReturnMatchPolicy{}).Where("ptype = ? AND v0 = ? AND v1 = ? AND v2 = ?", policy.Ptype, policy.V0, policy.V1, policy.V2).Count(&count).Error; err != nil {
return err
}
if count == 0 {
if err := tx.Create(&policy).Error; err != nil {
return err
}
}
}
return tx.Create(&common.Migration{Version: version}).Error
})
}
@@ -0,0 +1,63 @@
package version_local
import (
"testing"
"go-admin/app/goauto/access"
migrationmodels "go-admin/cmd/migrate/migration/models"
common "go-admin/common/models"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
func TestAfterSalesReturnMatchCatalogAndGrants(t *testing.T) {
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&migrationmodels.SysRole{}, &afterSalesReturnMatchAPI{}, &afterSalesReturnMatchPolicy{}, &common.Migration{}); err != nil {
t.Fatal(err)
}
if err = db.Create(&migrationmodels.SysRole{RoleKey: access.RoleAfterSales}).Error; err != nil {
t.Fatal(err)
}
for _, version := range []string{"test-first", "test-repeat"} {
if err = migrateAfterSalesReturnMatch(db, version); err != nil {
t.Fatal(err)
}
}
var apis []afterSalesReturnMatchAPI
if err = db.Find(&apis).Error; err != nil {
t.Fatal(err)
}
if len(apis) != 8 {
t.Fatalf("API count=%d", len(apis))
}
for _, api := range apis {
if api.Path == "" || api.Action == "" {
t.Fatalf("empty API path/action: id=%d", api.ID)
}
}
var count int64
if err = db.Model(&afterSalesReturnMatchPolicy{}).Where("v0 = ?", access.RoleAfterSales).Count(&count).Error; err != nil {
t.Fatal(err)
}
if count != 8 {
t.Fatalf("policy count=%d", count)
}
// Simulate the original deployed catalogue bug, then repair and repeat.
if err = db.Model(&afterSalesReturnMatchAPI{}).Where("id > 0").Updates(map[string]any{"path": "", "action": ""}).Error; err != nil {
t.Fatal(err)
}
for _, version := range []string{"repair-first", "repair-repeat"} {
if err = migrateFixReturnMatchAPICatalog(db, version); err != nil {
t.Fatal(err)
}
}
if err = db.Model(&afterSalesReturnMatchAPI{}).Where("path = '' OR action = ''").Count(&count).Error; err != nil {
t.Fatal(err)
}
if count != 0 {
t.Fatalf("unrepaired API rows=%d", count)
}
}
@@ -0,0 +1,34 @@
package version_local
import (
"runtime"
"go-admin/app/goauto/access"
"go-admin/cmd/migrate/migration"
common "go-admin/common/models"
"gorm.io/gorm"
)
// Repair the API catalogue rows created by 1789801200000 before the path and
// action fields were populated. This is idempotent and only touches the
// return-match endpoints.
func init() {
_, fileName, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateFixReturnMatchAPICatalog)
}
func migrateFixReturnMatchAPICatalog(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
for _, permission := range access.AdminAPIs {
if len(permission.Path) < len("/api/admin/v1/return-matches") || permission.Path[:len("/api/admin/v1/return-matches")] != "/api/admin/v1/return-matches" {
continue
}
if err := tx.Table("sys_api").Where("title = ?", permission.Title).Updates(map[string]any{
"title": permission.Title, "path": permission.Path, "action": permission.Method, "type": "BUS",
}).Error; err != nil {
return err
}
}
return tx.Create(&common.Migration{Version: version}).Error
})
}
@@ -0,0 +1,26 @@
package version_local
import (
"runtime"
"go-admin/app/goauto/models"
"go-admin/cmd/migrate/migration"
common "go-admin/common/models"
"gorm.io/gorm"
)
func init() {
_, file, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(file), migrateSYBSessionAuthLease)
}
// The lease is additive and contains no credentials. It serializes automatic
// SYB re-authentication across sync and purchase writeback workers.
func migrateSYBSessionAuthLease(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.AutoMigrate(&models.SYBSessionAuthLease{}); err != nil {
return err
}
return tx.Create(&common.Migration{Version: version}).Error
})
}