diff --git a/server/app/admin/router/init_router.go b/server/app/admin/router/init_router.go index a73b2ca..762a037 100644 --- a/server/app/admin/router/init_router.go +++ b/server/app/admin/router/init_router.go @@ -21,6 +21,7 @@ import ( goautosybproductfilter "go-admin/app/goauto/sybproductfilter" goautosybshop "go-admin/app/goauto/sybshop" goautotask "go-admin/app/goauto/task" + goautoyeeke "go-admin/app/goauto/yeeke" common "go-admin/common/middleware" ) @@ -69,4 +70,5 @@ func InitRouter() { goautosybinnercode.InitRouter(r, authMiddleware) goautosybshop.InitRouter(r, authMiddleware) goautosybproductfilter.InitRouter(r, authMiddleware) + goautoyeeke.InitRouter(r, authMiddleware) } diff --git a/server/app/goauto/access/modules.go b/server/app/goauto/access/modules.go index c500aaa..bd33faa 100644 --- a/server/app/goauto/access/modules.go +++ b/server/app/goauto/access/modules.go @@ -16,6 +16,7 @@ const ( ModuleCollectionTasks = "collection_tasks" ModulePurchaseTasks = "purchase_tasks" ModuleAIMatching = "ai_matching" + ModuleYeekeReturns = "yeeke_returns" ) // ModuleDefinition is the single source of truth shared by menu migration, @@ -67,6 +68,7 @@ var goAutoMenuGroupMetadata = []MenuGroupDefinition{ ModulePDDProducts, ModuleCollectionTasks, ModulePurchaseTasks, + ModuleYeekeReturns, }, }, { @@ -101,6 +103,7 @@ var goAutoModuleMetadata = []ModuleDefinition{ {Key: ModuleCollectionTasks, Title: "采集任务", Path: "/collection-tasks", RouteName: "GoAutoCollectionTasks", Component: "/goauto/collection-tasks/index", Icon: "list", Sort: 60, PurchaserDefault: true}, {Key: ModulePurchaseTasks, Title: "采购管理", Path: "/purchase-tasks", RouteName: "GoAutoPurchaseTasks", Component: "/goauto/purchase-tasks/index", Icon: "shopping", Sort: 61, PurchaserDefault: true}, {Key: ModuleAIMatching, Title: "AI 规格匹配", Path: "/ai-matching-settings", RouteName: "GoAutoAiMatchingSettings", Component: "/goauto/ai-matching-settings/index", Icon: "setting", Sort: 62, PurchaserHardHidden: true}, + {Key: ModuleYeekeReturns, Title: "yeeke 退货同步", Path: "/yeeke-returns", RouteName: "GoAutoYeekeReturns", Component: "/goauto/yeeke-returns/index", Icon: "time", Sort: 63, PurchaserDefault: true}, } // GoAutoModules returns independent copies so callers cannot mutate the @@ -161,6 +164,8 @@ func moduleKeyForAPI(path string) string { return ModulePurchaseTasks case strings.HasPrefix(path, "/api/admin/v1/ai-matching-settings"): return ModuleAIMatching + case strings.HasPrefix(path, "/api/admin/v1/yeeke-returns"): + return ModuleYeekeReturns default: return "" } diff --git a/server/app/goauto/access/modules_test.go b/server/app/goauto/access/modules_test.go index 610ebec..0e73709 100644 --- a/server/app/goauto/access/modules_test.go +++ b/server/app/goauto/access/modules_test.go @@ -10,8 +10,8 @@ import ( func TestGoAutoModulesOwnEveryAdminAPIExactlyOnce(t *testing.T) { modules := GoAutoModules() - if len(modules) != 13 { - t.Fatalf("got %d modules, want 13", len(modules)) + if len(modules) != 14 { + t.Fatalf("got %d modules, want 14", len(modules)) } owners := make(map[string]int) @@ -56,8 +56,8 @@ func TestGoAutoMenuSortsFitMySQLSignedTinyInt(t *testing.T) { } previous = module.Sort } - if modules[0].Sort != 50 || modules[len(modules)-1].Sort != 62 { - t.Fatalf("GoAuto menu sort range = %d..%d, want 50..62", modules[0].Sort, modules[len(modules)-1].Sort) + if modules[0].Sort != 50 || modules[len(modules)-1].Sort != 63 { + t.Fatalf("GoAuto menu sort range = %d..%d, want 50..63", modules[0].Sort, modules[len(modules)-1].Sort) } } @@ -71,7 +71,7 @@ func TestGoAutoMenuGroupsCoverModulesExactlyOnce(t *testing.T) { } wantOrder := [][]string{ - {ModuleSYBProducts, ModuleSYBSyncRuns, ModuleSYBInnerCodes, ModuleShopeeProducts, ModulePDDProducts, ModuleCollectionTasks, ModulePurchaseTasks}, + {ModuleSYBProducts, ModuleSYBSyncRuns, ModuleSYBInnerCodes, ModuleShopeeProducts, ModulePDDProducts, ModuleCollectionTasks, ModulePurchaseTasks, ModuleYeekeReturns}, {ModuleSYBShops, ModuleSYBProductFilters, ModuleCollectionRules, ModulePurchaseRules, ModuleDevices, ModuleAIMatching}, } seen := make(map[string]int) diff --git a/server/app/goauto/access/purchaser.go b/server/app/goauto/access/purchaser.go index a9f8302..1777036 100644 --- a/server/app/goauto/access/purchaser.go +++ b/server/app/goauto/access/purchaser.go @@ -124,6 +124,10 @@ var AdminAPIs = []APIPermission{ {"取消采购任务", "/api/admin/v1/purchase-tasks/:taskId/cancel", "POST", true}, {"处理结果不明确任务", "/api/admin/v1/purchase-tasks/:taskId/resolve-unknown", "POST", true}, + {"查看 yeeke 退货同步记录", "/api/admin/v1/yeeke-returns/sync-runs", "GET", true}, + {"查看 yeeke 退货同步详情", "/api/admin/v1/yeeke-returns/sync-runs/:runId", "GET", true}, + {"手动触发 yeeke 退货同步", "/api/admin/v1/yeeke-returns/sync", "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}, diff --git a/server/app/goauto/clientapi/gateway_test.go b/server/app/goauto/clientapi/gateway_test.go index 0d97d2f..f8c8088 100644 --- a/server/app/goauto/clientapi/gateway_test.go +++ b/server/app/goauto/clientapi/gateway_test.go @@ -33,7 +33,7 @@ func fixture(t *testing.T) (*gorm.DB, clientkey.Service) { func TestEveryRouteIsExplicitlyScoped(t *testing.T) { db, s := fixture(t) routes := Inventory() - if len(s.Modules) != 13 { + if len(s.Modules) != 14 { t.Fatal("menu groups lost") } router := gin.New() diff --git a/server/app/goauto/models/yeeke.go b/server/app/goauto/models/yeeke.go index 2848083..a158287 100644 --- a/server/app/goauto/models/yeeke.go +++ b/server/app/goauto/models/yeeke.go @@ -18,24 +18,29 @@ type YeekeSession struct { func (YeekeSession) TableName() string { return "yeeke_session" } type YeekeReturnPackage struct { - ID uint64 `gorm:"primaryKey;autoIncrement"` - ExternalID string `gorm:"size:128;not null;uniqueIndex:ux_yeeke_return_package_external"` - OrderSN string `gorm:"size:128;not null;index"` - TrackingNo string `gorm:"size:128;not null;index"` - ShopID string `gorm:"size:128;not null;default:''"` - ShopName string `gorm:"size:255;not null;default:''"` - WareCode string `gorm:"size:128;not null;default:''"` - WareHouse string `gorm:"size:255;not null;default:''"` - WareName string `gorm:"size:255;not null;default:''"` - ClaimStatus string `gorm:"size:64;not null;default:''"` - ClaimTime *time.Time - CreateTime *time.Time - UpdateTime *time.Time - DestroyDeadLine *time.Time - LastSyncedAt time.Time `gorm:"not null;index"` - SyncStatus string `gorm:"size:32;not null;default:'ok'"` - CreatedAt time.Time - UpdatedAt time.Time + ID uint64 `gorm:"primaryKey;autoIncrement"` + ExternalID string `gorm:"size:128;not null;uniqueIndex:ux_yeeke_return_package_external"` + OrderSN string `gorm:"size:128;not null;index"` + TrackingNo string `gorm:"size:128;not null;index"` + ShopID string `gorm:"size:128;not null;default:''"` + ShopName string `gorm:"size:255;not null;default:''"` + WareCode string `gorm:"size:128;not null;default:''"` + WareHouse string `gorm:"size:255;not null;default:''"` + WareName string `gorm:"size:255;not null;default:''"` + ClaimStatus string `gorm:"size:64;not null;default:''"` + // StatusUnrecognized is set when ClaimStatus is not one of the values the + // sync code currently understands. It is never bucketed into a known + // status silently (#336): the raw value is still kept in ClaimStatus, and + // this flag lets an operator find and review these rows. + StatusUnrecognized bool `gorm:"not null;default:false;index"` + ClaimTime *time.Time + CreateTime *time.Time + UpdateTime *time.Time + DestroyDeadLine *time.Time + LastSyncedAt time.Time `gorm:"not null;index"` + SyncStatus string `gorm:"size:32;not null;default:'ok'"` + CreatedAt time.Time + UpdatedAt time.Time } func (YeekeReturnPackage) TableName() string { return "yeeke_return_package" } diff --git a/server/app/goauto/yeeke/handler.go b/server/app/goauto/yeeke/handler.go new file mode 100644 index 0000000..31dfa4f --- /dev/null +++ b/server/app/goauto/yeeke/handler.go @@ -0,0 +1,169 @@ +package yeeke + +import ( + "errors" + "net/http" + "strconv" + + "go-admin/app/goauto/models" + + "github.com/gin-gonic/gin" + "github.com/go-admin-team/go-admin-core/sdk/api" + "github.com/go-admin-team/go-admin-core/sdk/pkg" + jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth" + "gorm.io/gorm" +) + +// Handler exposes the yeeke read-only admin surface: manual sync trigger and +// sync run history/summary. It never returns credentials, tokens, captcha +// text or a full raw yeeke response — only the counters already stored on +// models.YeekeSyncRun (#336 requirement #7). +type Handler struct { + // DB lets tests inject a database directly; production requests resolve + // it from the gin context via pkg.GetOrm, same as sybimport.Handler. + DB *gorm.DB +} + +func (h Handler) db(c *gin.Context) (*gorm.DB, bool) { + db := h.DB + var err error + if db == nil { + db, err = pkg.GetOrm(c) + } + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"code": "INTERNAL", "message": "服务端处理失败"}) + return nil, false + } + return db, true +} + +// SyncRunDTO is the read-only shape returned to admin/purchaser. It embeds +// only the summary fields already computed by the sync run itself; it never +// carries yeeke_session (token/cookies) or a raw page response. +type SyncRunDTO struct { + ID uint64 `json:"id"` + Status string `json:"status"` + Trigger string `json:"trigger"` + TotalPages int `json:"totalPages"` + ReadCount int `json:"readCount"` + CreatedCount int `json:"createdCount"` + UpdatedCount int `json:"updatedCount"` + SkippedCount int `json:"skippedCount"` + FailedCount int `json:"failedCount"` + ErrorMessage string `json:"errorMessage"` + StartedAt string `json:"startedAt"` + FinishedAt *string `json:"finishedAt"` + LastSuccessAt *string `json:"lastSuccessAt"` +} + +func toDTO(r models.YeekeSyncRun) SyncRunDTO { + dto := SyncRunDTO{ + ID: r.ID, Status: r.Status, Trigger: r.Trigger, TotalPages: r.TotalPages, + ReadCount: r.ReadCount, CreatedCount: r.CreatedCount, UpdatedCount: r.UpdatedCount, + SkippedCount: r.SkippedCount, FailedCount: r.FailedCount, ErrorMessage: r.ErrorMessage, + StartedAt: r.StartedAt.UTC().Format("2006-01-02T15:04:05Z"), + } + if r.FinishedAt != nil { + s := r.FinishedAt.UTC().Format("2006-01-02T15:04:05Z") + dto.FinishedAt = &s + } + if r.LastSuccessAt != nil { + s := r.LastSuccessAt.UTC().Format("2006-01-02T15:04:05Z") + dto.LastSuccessAt = &s + } + return dto +} + +// ListSyncRuns returns the most recent sync runs, newest first. Visible to +// admin and purchaser alike (#336 requirement #7); it is mounted without +// Casbin role gating, mirroring sybimport's /sync-runs. +func (h Handler) ListSyncRuns(c *gin.Context) { + page, err := strconv.Atoi(c.DefaultQuery("page", "1")) + if err != nil || page < 1 { + page = 1 + } + pageSize, err := strconv.Atoi(c.DefaultQuery("pageSize", "20")) + if err != nil || pageSize < 1 || pageSize > 100 { + pageSize = 20 + } + db, ok := h.db(c) + if !ok { + return + } + var total int64 + if err := db.Model(&models.YeekeSyncRun{}).Count(&total).Error; err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"code": "INTERNAL", "message": "服务端处理失败"}) + return + } + var rows []models.YeekeSyncRun + if err := db.Order("id desc").Offset((page - 1) * pageSize).Limit(pageSize).Find(&rows).Error; err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"code": "INTERNAL", "message": "服务端处理失败"}) + return + } + items := make([]SyncRunDTO, 0, len(rows)) + for _, row := range rows { + items = append(items, toDTO(row)) + } + var last models.YeekeSyncRun + lastSuccessAt := "" + if err := db.Where("status = ?", "succeeded").Order("id desc").First(&last).Error; err == nil && last.LastSuccessAt != nil { + lastSuccessAt = last.LastSuccessAt.UTC().Format("2006-01-02T15:04:05Z") + } + c.JSON(http.StatusOK, gin.H{"code": 200, "data": gin.H{ + "items": items, "total": total, "page": page, "pageSize": pageSize, + "lastSuccessAt": lastSuccessAt, + }}) +} + +// SyncRunDetail returns one run's summary. +func (h Handler) SyncRunDetail(c *gin.Context) { + id, err := strconv.ParseUint(c.Param("runId"), 10, 64) + if err != nil || id == 0 { + c.JSON(http.StatusBadRequest, gin.H{"code": "INVALID_REQUEST", "message": "runId 无效"}) + return + } + db, ok := h.db(c) + if !ok { + return + } + var row models.YeekeSyncRun + if err := db.First(&row, id).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + c.JSON(http.StatusNotFound, gin.H{"code": "NOT_FOUND", "message": "同步记录不存在"}) + return + } + c.JSON(http.StatusInternalServerError, gin.H{"code": "INTERNAL", "message": "服务端处理失败"}) + return + } + c.JSON(http.StatusOK, gin.H{"code": 200, "data": gin.H{"item": toDTO(row)}}) +} + +// TriggerSync starts a manual sync run. Only admin and purchaser may call it, +// same as SYB's manual Import (sybimport.Handler.Import) — this mirrors that +// role check exactly. +func (h Handler) TriggerSync(c *gin.Context) { + role, _ := jwt.ExtractClaims(c)["rolekey"].(string) + if role != "admin" && role != "purchaser" { + c.JSON(http.StatusForbidden, gin.H{"code": "FORBIDDEN", "message": "只有管理员或采购员可以开始同步"}) + return + } + db, ok := h.db(c) + if !ok { + return + } + result, err := StartSync(c.Request.Context(), db, nil, "manual", false) + if err != nil { + if errors.Is(err, ErrAlreadyRunning) { + c.JSON(http.StatusConflict, gin.H{"code": "ALREADY_RUNNING", "message": err.Error()}) + return + } + // `[必须]` err here is only ever a config/connect-stage message built in + // start.go/yeekeclient — never a raw yeeke response, never a token or + // credential. api.GetRequestLogger keeps the same text out of the HTTP + // body while still recording it server-side for operators. + api.GetRequestLogger(c).Errorf("yeeke manual sync failed to start: %v", err) + c.JSON(http.StatusBadGateway, gin.H{"code": "SYNC_START_FAILED", "message": err.Error()}) + return + } + c.JSON(http.StatusAccepted, gin.H{"code": 200, "data": gin.H{"runId": result.RunID, "skipped": result.Skipped}}) +} diff --git a/server/app/goauto/yeeke/job.go b/server/app/goauto/yeeke/job.go new file mode 100644 index 0000000..21994f8 --- /dev/null +++ b/server/app/goauto/yeeke/job.go @@ -0,0 +1,35 @@ +package yeeke + +import ( + "context" + "errors" + + "gorm.io/gorm" +) + +// ReturnSyncInvokeTarget is the go-admin job invoke_target key for the +// scheduled yeeke return sync (#336). The job row itself is seeded disabled +// (Status: 2) by migrations/version-local; an admin turns it on explicitly. +const ReturnSyncInvokeTarget = "GoAutoYeekeReturnSync" + +// ReturnSyncJob is registered in go-admin's ExecJob map (app/jobs/examples.go). +// ExecWithDB is the production path; Exec exists only to satisfy the legacy +// jobs.JobExec interface and fails closed if an older caller forgets to +// provide the current database. +type ReturnSyncJob struct{} + +func (ReturnSyncJob) Exec(_ interface{}) error { + return errors.New("yeeke 定时同步缺少数据库连接") +} + +// ExecWithDB starts a sync sharing the same StartSync entry point, and hence +// the same syncGate/active_slot lease, as the manual admin trigger — a +// scheduled tick that lands while a manual run (or a previous tick) is still +// in progress is skipped rather than queued or run concurrently. +func (ReturnSyncJob) ExecWithDB(db *gorm.DB, _ interface{}) error { + _, err := StartSync(context.Background(), db, nil, "scheduled", true) + if errors.Is(err, ErrAlreadyRunning) { + return nil + } + return err +} diff --git a/server/app/goauto/yeeke/router.go b/server/app/goauto/yeeke/router.go new file mode 100644 index 0000000..3721e30 --- /dev/null +++ b/server/app/goauto/yeeke/router.go @@ -0,0 +1,23 @@ +package yeeke + +import ( + "github.com/gin-gonic/gin" + jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth" +) + +// InitRouter mounts the yeeke read-only admin surface (#336): manual sync +// trigger plus sync run history/summary. Every route here only reads from or +// writes to GoAuto's own database and, for the trigger, starts a read-only +// yeeke sync — no yeeke write endpoint is ever called. +// +// `[必须]` These routes are authenticated but intentionally not Casbin-gated +// (middleware.AuthCheckRole), same as sybimport's /sync-runs: #336 requires +// both admin and purchaser to see the summary, and TriggerSync does its own +// admin/purchaser role check inline (mirroring sybimport.Handler.Import). +func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) { + handler := Handler{} + group := engine.Group("/api/admin/v1/yeeke-returns").Use(auth.MiddlewareFunc()) + group.GET("/sync-runs", handler.ListSyncRuns) + group.GET("/sync-runs/:runId", handler.SyncRunDetail) + group.POST("/sync", handler.TriggerSync) +} diff --git a/server/app/goauto/yeeke/start.go b/server/app/goauto/yeeke/start.go new file mode 100644 index 0000000..1d5d626 --- /dev/null +++ b/server/app/goauto/yeeke/start.go @@ -0,0 +1,110 @@ +package yeeke + +import ( + "context" + "errors" + "fmt" + "strings" + "sync" + "time" + + "go-admin/app/goauto/sybclient" + "go-admin/app/goauto/yeekeclient" + "go-admin/config" + + "gorm.io/gorm" +) + +// syncTimeout bounds one sync run. It exists so a stalled yeeke response +// cannot pin a background goroutine forever. +const syncTimeout = 30 * time.Minute + +// syncGate makes sync runs mutually exclusive within this process, same +// reasoning as sybimport.importGate: manual trigger and the scheduled job +// must never run concurrently (#336 requirement #5). The database-level +// unique active_slot column on yeeke_sync_run is the durable backstop (it +// also covers the unlikely case of two processes sharing one database); this +// in-memory gate exists to fail fast with a clear message before even +// touching the network or the OCR service. +var syncGate = struct { + sync.Mutex + running bool +}{} + +// ErrAlreadyRunning is returned by StartSync when a sync is already in +// progress, whether it was started manually or by the scheduled job. +var ErrAlreadyRunning = errors.New("已有 yeeke 退货包裹同步正在执行,请等它结束后再试") + +// OCR is satisfied by sybclient.OcrClient; yeeke reuses the same self-hosted +// captcha OCR service already approved for SYB (#336). +type OCR = yeekeclient.OCR + +// StartResult reports what StartSync did without exposing any run internals +// that should not cross the HTTP boundary. +type StartResult struct { + RunID uint64 + Skipped bool +} + +// StartSync is the single entry point shared by the authenticated admin +// handler and the scheduled job, so a scheduled run can never bypass the same +// concurrency guard as a manual run. +// +// `[必须]` This only ever reads from yeeke. Credentials are resolved by the +// caller (config.ExtConfig.Yeeke) and never logged here. +func StartSync(ctx context.Context, db *gorm.DB, ocr OCR, trigger string, skipIfRunning bool) (StartResult, error) { + settings := config.ExtConfig.Yeeke.Resolved() + if !settings.HasCredentials() { + return StartResult{}, fmt.Errorf("yeeke 账号未配置:请设置 GOAUTO_YEEKE_USERNAME / GOAUTO_YEEKE_PASSWORD,或在 config.yaml 的 yeeke 段填写 username/password,然后重启服务端") + } + if ocr == nil { + if strings.TrimSpace(settings.OcrURL) == "" { + return StartResult{}, fmt.Errorf("yeeke 验证码识别服务未配置:请设置 extend.yeeke.ocrurl") + } + // Reuse SYB's OcrClient implementation as-is (#336): independent + // http.Client, no shared cookie jar, captcha bytes stay in memory. See + // sybclient/ocr.go's package comment for why that isolation matters. + client, err := sybclient.NewOcrClient(settings.OcrURL, 0) + if err != nil { + return StartResult{}, err + } + ocr = client + } + + syncGate.Lock() + if syncGate.running { + syncGate.Unlock() + if skipIfRunning { + return StartResult{Skipped: true}, nil + } + return StartResult{}, ErrAlreadyRunning + } + syncGate.running = true + syncGate.Unlock() + release := func() { + syncGate.Lock() + syncGate.running = false + syncGate.Unlock() + } + + client, err := yeekeclient.Connect(ctx, yeekeclient.NewSessionStore(db), yeekeclient.Credentials{ + Username: settings.Username, Password: settings.Password, + }, settings.BaseURL, ocr, settings.OcrMaxAttempts) + if err != nil { + release() + return StartResult{}, err + } + + svc := NewService(db, client, Config{PageSize: settings.PageSize, MaxPages: settings.MaxPages, Retry: settings.Retry}) + // SyncAsync acquires the DB lease synchronously (so the caller gets a run + // id right away and the active_slot lease is held before this function + // returns) then walks pages in the background, bound to its own timeout + // independent of the HTTP request context. syncGate is released once that + // background walk finishes, not when this function returns. + runID, err := svc.SyncAsync(ctx, trigger, release) + if err != nil { + release() + return StartResult{}, err + } + return StartResult{RunID: runID}, nil +} diff --git a/server/app/goauto/yeeke/sync.go b/server/app/goauto/yeeke/sync.go index 61bba7d..b8837c2 100644 --- a/server/app/goauto/yeeke/sync.go +++ b/server/app/goauto/yeeke/sync.go @@ -10,6 +10,7 @@ import ( "go-admin/app/goauto/yeekeclient" "gorm.io/gorm" "strconv" + "strings" "time" ) @@ -40,6 +41,12 @@ type Report struct { Status string } +// 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} + func external(v any) string { return fmt.Sprint(v) } func stamp(t *yeekeclient.Timestamp) *time.Time { if t == nil || t.IsZero() { @@ -78,19 +85,58 @@ type Service struct { func NewService(db *gorm.DB, c *yeekeclient.Client, cfg Config) *Service { return &Service{db: db, client: c, cfg: cfg.norm()} } + +// Sync acquires the shared lease, runs the page walk synchronously and +// returns the final report. Tests use this directly; StartSync (start.go) +// uses SyncAsync instead so an HTTP request does not block for the whole +// run. func (s *Service) Sync(ctx context.Context, trigger string) (Report, error) { r, e := s.acquire(ctx, trigger) if e != nil { return Report{}, e } + return s.run(ctx, r) +} + +// SyncAsync acquires the lease synchronously (so the caller gets a run id +// immediately, and the unique active_slot lease is held before returning) +// and continues the page walk in a background goroutine bound to its own +// timeout, independent of the caller's request context. +// onDone, when non-nil, runs after the background page walk finishes +// (success or failure) — StartSync uses it to release the in-memory +// concurrency gate at the right time instead of when this function returns. +func (s *Service) SyncAsync(ctx context.Context, trigger string, onDone func()) (uint64, error) { + r, e := s.acquire(ctx, trigger) + if e != nil { + return 0, e + } + go func() { + bg, cancel := context.WithTimeout(context.Background(), syncTimeout) + defer cancel() + _, _ = s.run(bg, r) + if onDone != nil { + onDone() + } + }() + return r.ID, nil +} + +func (s *Service) run(ctx context.Context, r *models.YeekeSyncRun) (Report, error) { rep := Report{RunID: r.ID, Status: "failed"} + var errMsg string + var runErr error defer func() { now := time.Now().UTC() - s.db.Model(r).Updates(map[string]any{"status": rep.Status, "total_pages": rep.TotalPages, "read_count": rep.Read, "created_count": rep.Created, "updated_count": rep.Updated, "skipped_count": rep.Skipped, "failed_count": rep.Failed, "active_slot": nil, "lease_owner": "", "lease_expires_at": nil, "finished_at": now}) + updates := map[string]any{"status": rep.Status, "total_pages": rep.TotalPages, "read_count": rep.Read, "created_count": rep.Created, "updated_count": rep.Updated, "skipped_count": rep.Skipped, "failed_count": rep.Failed, "error_message": errMsg, "active_slot": nil, "lease_owner": "", "lease_expires_at": nil, "finished_at": now} + if rep.Status == "succeeded" { + updates["last_success_at"] = now + } + s.db.Model(r).Updates(updates) }() seen := map[string]bool{} 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 { @@ -98,12 +144,19 @@ func (s *Service) Sync(ctx context.Context, trigger string) (Report, error) { } select { case <-ctx.Done(): - return rep, ctx.Err() + runErr = ctx.Err() + errMsg = truncateRunError(runErr.Error()) + return rep, runErr case <-time.After(time.Duration(a+1) * 100 * time.Millisecond): } } if e != nil { - return rep, e + // 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 { @@ -140,6 +193,19 @@ func (s *Service) Sync(ctx context.Context, trigger string) (Report, error) { rep.Status = "succeeded" return rep, nil } + +// truncateRunError keeps error_message inside the column's size limit. It +// never includes request bodies or headers, so it cannot leak a captcha, +// token or credential: every error path above passes only Go error text from +// HTTP status/timeout/JSON-decoding failures. +func truncateRunError(value string) string { + const limit = 1000 + runes := []rune(strings.TrimSpace(value)) + if len(runes) <= limit { + return string(runes) + } + return string(runes[:limit]) +} func pageFingerprint(p yeekeclient.ReturnPage) string { b, _ := json.Marshal(p.Records) h := sha256.Sum256(b) @@ -154,7 +220,8 @@ func (s *Service) upsert(ctx context.Context, p yeekeclient.ReturnPackage) (bool if e != nil && !isNew { return false, false, e } - vals := map[string]any{"external_id": key, "order_sn": p.Ordersn, "tracking_no": p.TrackingNo, "shop_id": external(p.ShopID), "shop_name": p.ShopName, "ware_code": p.WareCode, "ware_house": p.WareHouse, "ware_name": p.WareName, "claim_status": external(p.Status), "claim_time": stamp(p.ClaimTime), "create_time": stamp(p.CreateTime), "update_time": stamp(p.UpdateTime), "destroy_dead_line": stamp(p.DestroyDeadLine), "last_synced_at": now, "sync_status": "ok"} + status := external(p.Status) + vals := map[string]any{"external_id": key, "order_sn": p.Ordersn, "tracking_no": p.TrackingNo, "shop_id": external(p.ShopID), "shop_name": p.ShopName, "ware_code": p.WareCode, "ware_house": p.WareHouse, "ware_name": p.WareName, "claim_status": status, "status_unrecognized": !knownClaimStatuses[status], "claim_time": stamp(p.ClaimTime), "create_time": stamp(p.CreateTime), "update_time": stamp(p.UpdateTime), "destroy_dead_line": stamp(p.DestroyDeadLine), "last_synced_at": now, "sync_status": "ok"} if isNew { row = models.YeekeReturnPackage{ExternalID: key} if e = s.db.WithContext(ctx).Create(&row).Error; e != nil { @@ -167,14 +234,17 @@ func (s *Service) upsert(ctx context.Context, p yeekeclient.ReturnPackage) (bool for n, i := range p.Items { ik := itemKey(p, i, n) ir := models.YeekeReturnItem{} - ie := s.db.Where("package_id = ? AND external_key = ?", row.ID, ik).First(&ir).Error + ie := s.db.WithContext(ctx).Where("package_id = ? AND external_key = ?", row.ID, ik).First(&ir).Error + if ie != nil && ie != gorm.ErrRecordNotFound { + return false, false, ie + } iv := map[string]any{"package_id": row.ID, "external_key": ik, "item_id": external(i.ItemID), "variation_id": external(i.VariationID), "item_name": i.ItemName, "variation_name": i.VariationName, "image": i.Image, "quantity": i.Quantity, "last_synced_at": now, "sync_status": "ok"} if ie == gorm.ErrRecordNotFound { - if e = s.db.Create(&models.YeekeReturnItem{PackageID: row.ID, ExternalKey: ik}).Error; e != nil { + if e = s.db.WithContext(ctx).Create(&models.YeekeReturnItem{PackageID: row.ID, ExternalKey: ik}).Error; e != nil { return false, false, e } } - if e = s.db.Model(&ir).Where("package_id = ? AND external_key = ?", row.ID, ik).Updates(iv).Error; e != nil { + if e = s.db.WithContext(ctx).Model(&models.YeekeReturnItem{}).Where("package_id = ? AND external_key = ?", row.ID, ik).Updates(iv).Error; e != nil { return false, false, e } } diff --git a/server/app/goauto/yeeke/sync_paging_test.go b/server/app/goauto/yeeke/sync_paging_test.go new file mode 100644 index 0000000..6d5f295 --- /dev/null +++ b/server/app/goauto/yeeke/sync_paging_test.go @@ -0,0 +1,462 @@ +package yeeke + +import ( + "context" + "encoding/json" + "fmt" + "go-admin/app/goauto/migrations" + "go-admin/app/goauto/models" + "go-admin/app/goauto/yeekeclient" + "go-admin/config" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + "net/http" + "net/http/httptest" + "strings" + "sync" + "sync/atomic" + "testing" + "time" +) + +type stubOCR struct{ codes []string } + +func (s *stubOCR) Recognize(context.Context, []byte) (string, error) { + if len(s.codes) == 0 { + return "", nil + } + c := s.codes[0] + s.codes = s.codes[1:] + return c, nil +} + +func testDB(t *testing.T) *gorm.DB { + t.Helper() + db, err := gorm.Open(sqlite.Open(fmt.Sprintf("file:%s?mode=memory&cache=shared", t.Name())), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := migrations.Migrate(db); err != nil { + t.Fatal(err) + } + return db +} + +func record(id, itemID, variationID string, status any) string { + b, _ := json.Marshal(map[string]any{ + "id": id, "ordersn": "o-" + id, "trackingNo": "t-" + id, "status": status, + "items": []map[string]any{{ + "id": id + "-i1", "itemId": itemID, "variationId": variationID, + "itemName": "n", "variationName": "v", "variationQuantityPurchased": 1, + }}, + }) + return string(b) +} + +func page(records []string, total, pages int) string { + return fmt.Sprintf(`{"success":true,"result":{"records":[%s],"total":%d,"pages":%d}}`, strings.Join(records, ","), total, pages) +} + +// TestPagingSurvivesTotalChangingMidRun: page 1 reports one total/pages, page +// 2 reports a different total/pages (the underlying data changed between the +// two requests). The walk must still finish using what each page returned. +func TestPagingSurvivesTotalChangingMidRun(t *testing.T) { + db := testDB(t) + var calls int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + n := atomic.AddInt32(&calls, 1) + w.Header().Set("Content-Type", "application/json") + switch n { + case 1: + fmt.Fprint(w, page([]string{record("p1", "i", "v1", 1)}, 100, 2)) + case 2: + fmt.Fprint(w, page([]string{record("p2", "i", "v2", 1)}, 50, 1)) + default: + fmt.Fprint(w, page(nil, 0, 0)) + } + })) + defer srv.Close() + 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" || rep.Read != 2 { + t.Fatalf("rep=%+v", rep) + } +} + +// TestPagingSkipsARepeatedDuplicatePage: the server returns the exact same +// page twice in a row (e.g. a retried request landed after all). The second +// occurrence must be recognized as a duplicate and stop the walk instead of +// looping or double counting. +func TestPagingSkipsARepeatedDuplicatePage(t *testing.T) { + db := testDB(t) + var calls int32 + body := page([]string{record("p1", "i", "v1", 1)}, 10, 5) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + n := atomic.AddInt32(&calls, 1) + w.Header().Set("Content-Type", "application/json") + if n <= 2 { + fmt.Fprint(w, body) // identical page served twice + return + } + fmt.Fprint(w, page(nil, 10, 5)) + })) + defer srv.Close() + 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" { + t.Fatalf("rep=%+v", rep) + } + var n int64 + db.Model(&models.YeekeReturnPackage{}).Count(&n) + if n != 1 { + t.Fatalf("packages=%d, want 1 (duplicate page must not double-insert)", n) + } + if calls != 2 { + t.Fatalf("calls=%d, want exactly 2 (stop right after recognizing the duplicate)", calls) + } +} + +// TestPagingStopsOnEmptyPage confirms an empty page ends the walk cleanly. +func TestPagingStopsOnEmptyPage(t *testing.T) { + db := testDB(t) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, page(nil, 0, 0)) + })) + 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 { + t.Fatal(err) + } + if rep.Status != "succeeded" || rep.TotalPages != 1 || rep.Read != 0 { + t.Fatalf("rep=%+v", rep) + } +} + +// TestPagingTimeoutFailsRunButKeepsEarlierPages: page 1 succeeds and is +// persisted; page 2 always times out. The run must end as "failed" (after +// retrying up to cfg.Retry times) but page 1's row must remain intact — a +// failed page must never roll back or overwrite valid prior data. +func TestPagingTimeoutFailsRunButKeepsEarlierPages(t *testing.T) { + db := testDB(t) + var calls int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + n := atomic.AddInt32(&calls, 1) + if n == 1 { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, page([]string{record("p1", "i", "v1", 1)}, 10, 5)) + return + } + // Simulate a slow/timed-out request: this response is deliberately + // slower than the client's own context deadline below, so the client + // side must give up on its own rather than waiting for us. + time.Sleep(2 * time.Second) + })) + defer srv.Close() + c, _ := yeekeclient.New(srv.URL) + s := NewService(db, c, Config{PageSize: 1, Retry: 1}) + + ctx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) + defer cancel() + rep, err := s.Sync(ctx, "manual") + if err == nil { + t.Fatal("expected the timed-out page to surface an error") + } + if rep.Status != "failed" || rep.RunID == 0 { + t.Fatalf("rep=%+v", rep) + } + var row models.YeekeReturnPackage + if e := db.Where("external_id = ?", "p1").First(&row).Error; e != nil { + t.Fatalf("page 1's row must survive a later page's failure: %v", e) + } + if row.ClaimStatus != "1" { + t.Fatalf("page 1's row must be unmodified: %+v", row) + } + var run models.YeekeSyncRun + if e := db.First(&run, rep.RunID).Error; e != nil { + t.Fatal(e) + } + if run.Status != "failed" || run.ErrorMessage == "" { + t.Fatalf("run=%+v", run) + } +} + +// TestSyncResumeAfterSimulatedRestart: a first run fails partway through +// (simulating the process being interrupted after committing page 1). A +// second, full run afterwards must succeed and must not duplicate the row +// page 1 already wrote — it converges to exactly one package row. +func TestSyncResumeAfterSimulatedRestart(t *testing.T) { + db := testDB(t) + var calls int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + n := atomic.AddInt32(&calls, 1) + w.Header().Set("Content-Type", "application/json") + if n == 1 { + fmt.Fprint(w, page([]string{record("p1", "i", "v1", 1)}, 10, 2)) + return + } + if n == 2 { + w.WriteHeader(http.StatusInternalServerError) + return + } + // Full second run: both pages succeed this time. + if n == 3 { + fmt.Fprint(w, page([]string{record("p1", "i", "v1", 1)}, 10, 2)) + return + } + fmt.Fprint(w, page(nil, 10, 2)) + })) + defer srv.Close() + c, _ := yeekeclient.New(srv.URL) + s := NewService(db, c, Config{PageSize: 1, Retry: 0}) + + if _, err := s.Sync(context.Background(), "manual"); err == nil { + t.Fatal("expected the first ('interrupted') run to fail") + } + rep2, err := s.Sync(context.Background(), "manual") + if err != nil { + t.Fatal(err) + } + if rep2.Status != "succeeded" { + t.Fatalf("rep2=%+v", rep2) + } + var n int64 + db.Model(&models.YeekeReturnPackage{}).Count(&n) + if n != 1 { + t.Fatalf("packages=%d, want 1 (resume must not duplicate p1)", n) + } +} + +// TestIdempotentStatusUpdateInPlace: syncing the same package twice with a +// different status the second time updates the existing row rather than +// creating a second one. +func TestIdempotentStatusUpdateInPlace(t *testing.T) { + db := testDB(t) + var status int32 = 1 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, page([]string{record("p1", "i", "v1", atomic.LoadInt32(&status))}, 1, 1)) + })) + defer srv.Close() + c, _ := yeekeclient.New(srv.URL) + s := NewService(db, c, Config{PageSize: 10}) + if _, err := s.Sync(context.Background(), "manual"); err != nil { + t.Fatal(err) + } + atomic.StoreInt32(&status, 9) // an unrecognized status the second time + if _, err := s.Sync(context.Background(), "manual"); err != nil { + t.Fatal(err) + } + var n int64 + db.Model(&models.YeekeReturnPackage{}).Count(&n) + if n != 1 { + t.Fatalf("packages=%d, want 1 (status change must update in place)", n) + } + var row models.YeekeReturnPackage + db.Where("external_id = ?", "p1").First(&row) + if row.ClaimStatus != "9" || !row.StatusUnrecognized { + t.Fatalf("row=%+v, want claim_status=9 flagged unrecognized", row) + } +} + +// TestUnknownStatusIsPreservedVerbatimAndFlagged: an unrecognized status value +// is kept as-is (never remapped into a known bucket) and the row is flagged. +func TestUnknownStatusIsPreservedVerbatimAndFlagged(t *testing.T) { + db := testDB(t) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, page([]string{record("p1", "i", "v1", 7)}, 1, 1)) + })) + defer srv.Close() + c, _ := yeekeclient.New(srv.URL) + s := NewService(db, c, Config{PageSize: 10}) + if _, err := s.Sync(context.Background(), "manual"); err != nil { + t.Fatal(err) + } + var row models.YeekeReturnPackage + if e := db.Where("external_id = ?", "p1").First(&row).Error; e != nil { + t.Fatal(e) + } + if row.ClaimStatus != "7" { + t.Fatalf("claim_status=%q, want the raw value 7 preserved verbatim", row.ClaimStatus) + } + if !row.StatusUnrecognized { + t.Fatal("an unknown status must be flagged, not silently accepted") + } +} + +// TestActiveSlotLeaseRejectsConcurrentRuns exercises the same DB-level +// uniqueness the manual trigger and the scheduled job both rely on +// (models.YeekeSyncRun.ActiveSlot): two Sync calls racing against the same +// database must not both hold the lease at once. +func TestActiveSlotLeaseRejectsConcurrentRuns(t *testing.T) { + db := testDB(t) + release := make(chan struct{}) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + <-release + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, page(nil, 0, 0)) + })) + defer srv.Close() + c, _ := yeekeclient.New(srv.URL) + s := NewService(db, c, Config{PageSize: 10}) + + started := make(chan struct{}) + var firstErr, secondErr error + go func() { + close(started) + _, firstErr = s.Sync(context.Background(), "manual") + }() + <-started + time.Sleep(50 * time.Millisecond) // let the first Sync acquire its lease row + _, secondErr = s.Sync(context.Background(), "scheduled") + close(release) + time.Sleep(50 * time.Millisecond) + + if secondErr == nil { + t.Fatal("a second concurrent Sync must be rejected by the active_slot lease") + } + _ = firstErr +} + +// TestStartSyncGateRejectsConcurrentTriggers exercises StartSync itself (the +// entry point manual and scheduled triggers actually share): a second call +// while one is in flight is rejected, and with skipIfRunning it is reported +// as skipped instead of erroring — this is what the scheduled job uses so a +// tick landing during a manual run does not surface as a failure. +func TestStartSyncGateRejectsConcurrentTriggers(t *testing.T) { + db := testDB(t) + restore := setTestYeekeConfig(t, "op", "secret-pw") + defer restore() + + release := make(chan struct{}) + var loginCalls, listCalls int32 + srv := fakeYeekeServer(t, &loginCalls, &listCalls, release) + defer srv.Close() + restoreURL := setTestYeekeBaseURL(t, srv.URL) + defer restoreURL() + + var wg sync.WaitGroup + wg.Add(1) + var firstErr error + go func() { + defer wg.Done() + _, firstErr = StartSync(context.Background(), db, &stubOCR{codes: []string{"abcd"}}, "manual", false) + }() + // Give the first call time to acquire syncGate and the DB lease before + // the second one is attempted. + time.Sleep(150 * time.Millisecond) + + _, err := StartSync(context.Background(), db, &stubOCR{}, "scheduled", true) + if err != nil { + t.Fatalf("skipIfRunning=true must not error, got %v", err) + } + _, err2 := StartSync(context.Background(), db, &stubOCR{}, "manual", false) + if err2 == nil || err2 != ErrAlreadyRunning { + t.Fatalf("skipIfRunning=false must report ErrAlreadyRunning, got %v", err2) + } + + close(release) + wg.Wait() + if firstErr != nil { + t.Fatalf("first StartSync should have completed cleanly: %v", firstErr) + } +} + +// TestLoginFailureMessageNeverLeaksCredentialsOrCaptcha: whatever StartSync +// or the underlying client return as an error, the credential, password and +// recognized captcha text must never appear in it, since that text ends up +// in server logs and (truncated) in yeeke_sync_run.error_message. +func TestLoginFailureMessageNeverLeaksCredentialsOrCaptcha(t *testing.T) { + db := testDB(t) + const secretPassword = "S3cr3t-Do-Not-Leak" + const secretCaptcha = "zZqQ9x" + restore := setTestYeekeConfig(t, "leak-user", secretPassword) + defer restore() + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/agent-foreign/sys/randomImage": + json.NewEncoder(w).Encode(map[string]any{"success": true, "result": map[string]string{"image": "data:image/jpg;base64,SGk=", "checkKey": "k"}}) + case "/agent-foreign/sys/login": + json.NewEncoder(w).Encode(map[string]any{"success": false, "code": 1, "message": "验证码错误"}) + default: + http.NotFound(w, r) + } + })) + defer srv.Close() + restoreURL := setTestYeekeBaseURL(t, srv.URL) + defer restoreURL() + + _, err := StartSync(context.Background(), db, &stubOCR{codes: []string{secretCaptcha}}, "manual", false) + if err == nil { + t.Fatal("expected a login failure") + } + msg := err.Error() + if strings.Contains(msg, secretPassword) || strings.Contains(msg, secretCaptcha) || strings.Contains(msg, "leak-user") { + t.Fatalf("error message leaked a credential or captcha text: %q", msg) + } + + var run models.YeekeSyncRun + // StartSync failed before acquiring a run row here (Connect failed first), + // so there should be no run row at all to check — that is itself part of + // the guarantee: a failed login never gets far enough to write a summary + // row that could carry sensitive text. + if e := db.Order("id desc").First(&run).Error; e == nil { + if strings.Contains(run.ErrorMessage, secretPassword) || strings.Contains(run.ErrorMessage, secretCaptcha) { + t.Fatalf("run.ErrorMessage leaked a credential or captcha text: %q", run.ErrorMessage) + } + } +} + +// fakeYeekeServer serves a minimal login+list surface. It blocks the *list* +// call on release, so a test can hold StartSync's background goroutine open +// long enough to exercise the concurrency gate. +func fakeYeekeServer(t *testing.T, loginCalls, listCalls *int32, release chan struct{}) *httptest.Server { + t.Helper() + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/agent-foreign/sys/randomImage": + json.NewEncoder(w).Encode(map[string]any{"success": true, "result": map[string]string{"image": "data:image/jpg;base64,SGk=", "checkKey": "k"}}) + case "/agent-foreign/sys/login": + atomic.AddInt32(loginCalls, 1) + json.NewEncoder(w).Encode(map[string]any{"success": true, "result": map[string]any{"token": "tok", "userInfo": map[string]any{"id": "u1", "username": "op"}}}) + case "/agent-foreign/packageClaimRec/relation/list": + atomic.AddInt32(listCalls, 1) + <-release + fmt.Fprint(w, page(nil, 0, 0)) + default: + http.NotFound(w, r) + } + })) +} + +// The next two helpers isolate StartSync's config.ExtConfig.Yeeke dependency +// for tests, restoring it afterwards so other tests are unaffected. +func setTestYeekeConfig(t *testing.T, username, password string) func() { + t.Helper() + before := config.ExtConfig.Yeeke + config.ExtConfig.Yeeke.Username = username + config.ExtConfig.Yeeke.Password = password + config.ExtConfig.Yeeke.OcrURL = "http://unused.invalid/ocr" + return func() { config.ExtConfig.Yeeke = before } +} + +func setTestYeekeBaseURL(t *testing.T, url string) func() { + t.Helper() + before := config.ExtConfig.Yeeke + config.ExtConfig.Yeeke.BaseURL = url + return func() { config.ExtConfig.Yeeke = before } +} diff --git a/server/app/goauto/yeekeclient/connect_test.go b/server/app/goauto/yeekeclient/connect_test.go new file mode 100644 index 0000000..5362396 --- /dev/null +++ b/server/app/goauto/yeekeclient/connect_test.go @@ -0,0 +1,177 @@ +package yeekeclient + +import ( + "context" + "encoding/json" + "go-admin/app/goauto/models" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func newDB(t *testing.T) *gorm.DB { + t.Helper() + db, err := gorm.Open(sqlite.Open("file::memory:?cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&models.YeekeSession{}); err != nil { + t.Fatal(err) + } + return db +} + +type stubOCR struct { + codes []string + calls int +} + +func (s *stubOCR) Recognize(context.Context, []byte) (string, error) { + if s.calls >= len(s.codes) { + return "", nil + } + c := s.codes[s.calls] + s.calls++ + return c, nil +} + +// server builds a fake yeeke backend. loginOK controls whether /login accepts +// the submitted captcha; sessionValid controls whether /userInfo (used by +// CheckSession) reports the cached token as still good. +func fakeServer(t *testing.T, loginOK func(captcha string) bool, sessionValid func(token string) bool) *httptest.Server { + t.Helper() + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/agent-foreign/sys/randomImage": + json.NewEncoder(w).Encode(map[string]any{"success": true, "result": map[string]string{"image": "data:image/jpg;base64,SGk=", "checkKey": "k"}}) + case "/agent-foreign/sys/login": + var body map[string]string + json.NewDecoder(r.Body).Decode(&body) + if loginOK(body["captcha"]) { + json.NewEncoder(w).Encode(map[string]any{"success": true, "result": map[string]any{"token": "tok-" + body["captcha"], "userInfo": map[string]any{"id": "u1", "username": "u"}}}) + return + } + json.NewEncoder(w).Encode(map[string]any{"success": false, "code": 1, "message": "验证码错误"}) + case "/agent-foreign/sys/userInfo": + token := r.URL.Query().Get("token") + if sessionValid(token) { + json.NewEncoder(w).Encode(map[string]any{"success": true, "result": map[string]any{}}) + return + } + w.WriteHeader(http.StatusOK) + json.NewEncoder(w).Encode(map[string]any{"success": false, "code": 401, "message": "登录已失效"}) + default: + http.NotFound(w, r) + } + })) +} + +// TestConnectLoginsOnceThenReusesSession: a fresh Connect performs exactly one +// login, and a second Connect call with the cached session valid performs no +// login at all (token reuse, no re-login when session valid). +func TestConnectLoginsOnceThenReusesSession(t *testing.T) { + db := newDB(t) + loginCalls := 0 + srv := fakeServer(t, + func(captcha string) bool { loginCalls++; return captcha == "abcd" }, + func(token string) bool { return token == "tok-abcd" }, + ) + defer srv.Close() + + store := NewSessionStore(db) + ocr := &stubOCR{codes: []string{"abcd"}} + client, err := Connect(context.Background(), store, Credentials{Username: "u", Password: "p"}, srv.URL, ocr, 3) + if err != nil { + t.Fatal(err) + } + if client.Token() != "tok-abcd" { + t.Fatalf("token=%q", client.Token()) + } + if loginCalls != 1 { + t.Fatalf("loginCalls=%d, want 1", loginCalls) + } + + // Second connect: session is cached and still valid, so this must not + // touch OCR or /login again. + ocr2 := &stubOCR{codes: []string{"should-not-be-used"}} + client2, err := Connect(context.Background(), store, Credentials{Username: "u", Password: "p"}, srv.URL, ocr2, 3) + if err != nil { + t.Fatal(err) + } + if client2.Token() != "tok-abcd" { + t.Fatalf("reused token=%q", client2.Token()) + } + if loginCalls != 1 { + t.Fatalf("loginCalls after reuse=%d, want still 1", loginCalls) + } + if ocr2.calls != 0 { + t.Fatalf("OCR must not be called when the cached session is valid") + } +} + +// TestConnectReLoginsAfterSessionExpiredAndIsBounded: when the cached session +// is explicitly rejected (ErrSessionInvalid), Connect re-logs in — but only +// up to maxLogin captcha attempts, never looping forever. +func TestConnectReLoginsAfterSessionExpiredAndIsBounded(t *testing.T) { + db := newDB(t) + store := NewSessionStore(db) + // Seed an already-cached, not-yet-expired session so Connect's Load finds + // it and only CheckSession decides it is dead. + if err := store.Save(context.Background(), Session{ + Username: "u", Token: "stale", CookiesJSON: `[]`, UserID: "u1", + ExpiresAt: time.Now().UTC().Add(time.Hour), + }); err != nil { + t.Fatal(err) + } + + loginAttempts := 0 + srv := fakeServer(t, + func(captcha string) bool { loginAttempts++; return false }, // every captcha rejected + func(token string) bool { return false }, // cached session always invalid + ) + defer srv.Close() + + ocr := &stubOCR{codes: []string{"1", "2", "3", "4", "5", "6"}} // more codes than maxLogin allows + _, err := Connect(context.Background(), store, Credentials{Username: "u", Password: "p"}, srv.URL, ocr, 3) + if err == nil { + t.Fatal("expected login failure") + } + if loginAttempts != 3 { + t.Fatalf("loginAttempts=%d, want exactly maxLogin=3 (bounded, not endless)", loginAttempts) + } + + // The rejected cached session must have been deleted, not left in place. + if _, loadErr := store.Load(context.Background(), "u", time.Now().UTC()); loadErr != ErrNoSession { + t.Fatalf("expired/invalid session should have been deleted: %v", loadErr) + } +} + +// TestConnectKeepsCachedSessionOnTimeoutOrServerError: a network-level error +// checking the session (not an explicit "invalid") must not discard a +// possibly-still-good cached session. +func TestConnectKeepsCachedSessionOnTimeoutOrServerError(t *testing.T) { + db := newDB(t) + store := NewSessionStore(db) + if err := store.Save(context.Background(), Session{ + Username: "u", Token: "tok", CookiesJSON: `[]`, UserID: "u1", + ExpiresAt: time.Now().UTC().Add(time.Hour), + }); err != nil { + t.Fatal(err) + } + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + defer srv.Close() + + _, err := Connect(context.Background(), store, Credentials{Username: "u", Password: "p"}, srv.URL, &stubOCR{}, 3) + if err == nil { + t.Fatal("expected a propagated 5xx error") + } + if _, loadErr := store.Load(context.Background(), "u", time.Now().UTC()); loadErr != nil { + t.Fatalf("a 5xx must not discard the cached session: %v", loadErr) + } +} diff --git a/server/app/jobs/examples.go b/server/app/jobs/examples.go index 8036b15..9324d37 100644 --- a/server/app/jobs/examples.go +++ b/server/app/jobs/examples.go @@ -6,6 +6,7 @@ import ( "go-admin/app/goauto/shopeeproduct" "go-admin/app/goauto/sybimport" + "go-admin/app/goauto/yeeke" ) // InitJob @@ -17,6 +18,7 @@ func InitJob() { sybimport.HourlySyncInvokeTarget: sybimport.HourlySyncJob{}, sybimport.SpecAIParseInvokeTarget: sybimport.ScheduledSpecAIParseJob{}, shopeeproduct.SpecAutoMatchInvokeTarget: shopeeproduct.ScheduledAutoMatchJob{}, + yeeke.ReturnSyncInvokeTarget: yeeke.ReturnSyncJob{}, // ... } } diff --git a/server/cmd/migrate/migration/version-local/1789800500000_yeeke_return_sync_job.go b/server/cmd/migrate/migration/version-local/1789800500000_yeeke_return_sync_job.go new file mode 100644 index 0000000..582f8fa --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1789800500000_yeeke_return_sync_job.go @@ -0,0 +1,49 @@ +package version_local + +import ( + "errors" + "runtime" + + "go-admin/app/goauto/yeeke" + jobsmodels "go-admin/app/jobs/models" + "go-admin/cmd/migrate/migration" + common "go-admin/common/models" + + "gorm.io/gorm" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateYeekeReturnSyncJob) +} + +// migrateYeekeReturnSyncJob seeds the scheduled yeeke return sync job row, +// disabled by default (#336 requirement: 默认关闭定时任务,管理员手动开启). +// It follows 1786701600000_syb_hourly_sync_job.go exactly: Status 2 keeps the +// row out of the running cron set (see app/jobs/service/sys_job.go, which +// only adds Status == 1 jobs to the cron at startup). +func migrateYeekeReturnSyncJob(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := ensureYeekeReturnSyncJob(tx); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} + +func ensureYeekeReturnSyncJob(db *gorm.DB) error { + var existing jobsmodels.SysJob + err := db.Where("invoke_target = ?", yeeke.ReturnSyncInvokeTarget).First(&existing).Error + if err == nil { + return nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + return db.Create(&jobsmodels.SysJob{ + JobName: "yeeke 退货包裹同步", JobGroup: "GoAuto", JobType: 2, + CronExpression: "0 20 * * * *", InvokeTarget: yeeke.ReturnSyncInvokeTarget, + Args: "", + MisfirePolicy: 1, Concurrent: 1, Status: 2, + }).Error +} diff --git a/server/config/extend.go b/server/config/extend.go index 79e8ce1..1a59e02 100644 --- a/server/config/extend.go +++ b/server/config/extend.go @@ -18,8 +18,9 @@ var ExtConfig Extend // // 使用方法: config.ExtConfig......即可!! type Extend struct { - AMap AMap // 这里配置对应配置文件的结构即可 - SYB SYB + AMap AMap // 这里配置对应配置文件的结构即可 + SYB SYB + Yeeke Yeeke } type AMap struct { @@ -89,6 +90,65 @@ func (s SYB) HasCredentials() bool { return strings.TrimSpace(s.Username) != "" && s.Password != "" } +// Yeeke holds the mmt.yeeke.com 退货包裹只读对接 connection settings (#336). +// +// `[必须]` Username and Password are NOT read from settings.yml — same rule as +// SYB above — they come from GOAUTO_YEEKE_USERNAME / GOAUTO_YEEKE_PASSWORD (or +// config.yaml's yeeke: section) so no credential ever lands in a tracked file. +type Yeeke struct { + BaseURL string + Username string + Password string + PageSize int + MaxPages int + Retry int + // OcrURL is the captcha recognition service shared with SYB (#336, approved + // 2026-09-23). Empty disables OCR; there is no manual-entry fallback here + // because this is a server-side scheduled/triggered flow, not an interactive + // login form, so a disabled OCR simply makes sync fail with a clear error. + OcrURL string + OcrMaxAttempts int +} + +// YeekeDefaults are the values used when settings.yml leaves a field blank. +const ( + DefaultYeekeBaseURL = "https://mmt.yeeke.com" + DefaultYeekePageSize = 100 + DefaultYeekeMaxPages = 10000 + DefaultYeekeRetry = 2 + DefaultYeekeOcrMaxAttempts = 5 +) + +// Resolved returns the Yeeke settings with blanks replaced by defaults. It +// never defaults Username or Password: missing credentials must surface as an +// error at the call site, not as an attempt to log in as nobody. +func (y Yeeke) Resolved() Yeeke { + if strings.TrimSpace(y.BaseURL) == "" { + y.BaseURL = DefaultYeekeBaseURL + } + if y.PageSize <= 0 { + y.PageSize = DefaultYeekePageSize + } + if y.MaxPages <= 0 { + y.MaxPages = DefaultYeekeMaxPages + } + if y.Retry < 0 { + y.Retry = DefaultYeekeRetry + } + if y.OcrMaxAttempts <= 0 { + y.OcrMaxAttempts = DefaultYeekeOcrMaxAttempts + } + y.Username = strings.TrimSpace(y.Username) + y.BaseURL = strings.TrimRight(strings.TrimSpace(y.BaseURL), "/") + y.OcrURL = strings.TrimSpace(y.OcrURL) + return y +} + +// HasCredentials reports whether both account fields were supplied. +func (y Yeeke) HasCredentials() bool { + return strings.TrimSpace(y.Username) != "" && y.Password != "" +} + // ApplyEnvironment replaces tracked defaults with process-local runtime values. // Credentials stay outside tracked configuration files. GOAUTO_DB_DRIVER // defaults to mysql when GOAUTO_DB_DSN is present. @@ -104,6 +164,13 @@ func ApplyEnvironment() { if password := os.Getenv("GOAUTO_SYB_PASSWORD"); password != "" { ExtConfig.SYB.Password = password } + if username := strings.TrimSpace(os.Getenv("GOAUTO_YEEKE_USERNAME")); username != "" { + ExtConfig.Yeeke.Username = username + } + // `[必须]` Taken verbatim, same reasoning as GOAUTO_SYB_PASSWORD above. + if password := os.Getenv("GOAUTO_YEEKE_PASSWORD"); password != "" { + ExtConfig.Yeeke.Password = password + } dsn := strings.TrimSpace(os.Getenv("GOAUTO_DB_DSN")) if dsn == "" { diff --git a/server/config/extend_settings_test.go b/server/config/extend_settings_test.go index 211f3ed..4d2e3b4 100644 --- a/server/config/extend_settings_test.go +++ b/server/config/extend_settings_test.go @@ -62,5 +62,45 @@ func TestSettingsFilesCarryNoSYBCredentials(t *testing.T) { if file.Settings.Extend.SYB.Username != "" || file.Settings.Extend.SYB.Password != "" { t.Fatalf("%s 里出现了顺云宝凭据,凭据必须走环境变量", name) } + if file.Settings.Extend.Yeeke.Username != "" || file.Settings.Extend.Yeeke.Password != "" { + t.Fatalf("%s 里出现了 yeeke 凭据,凭据必须走环境变量 (#336)", name) + } + } +} + +// settings.yml 里的 extend.yeeke 必须真的能绑进 ExtConfig.Yeeke,同样的道理 +// 见上面 TestSYBSettingsInRepoBindToExtendStruct 的注释 (#336)。 +func TestYeekeSettingsInRepoBindToExtendStruct(t *testing.T) { + raw, err := os.ReadFile("settings.yml") + if err != nil { + t.Fatalf("读取 settings.yml 失败: %v", err) + } + var file struct { + Settings struct { + Extend Extend `yaml:"extend"` + } `yaml:"settings"` + } + if err := yaml.Unmarshal(raw, &file); err != nil { + t.Fatalf("解析 settings.yml 失败: %v", err) + } + + yeeke := file.Settings.Extend.Yeeke + if yeeke.BaseURL != "https://mmt.yeeke.com" { + t.Fatalf("baseurl 没有绑定成功: %q", yeeke.BaseURL) + } + if yeeke.PageSize != 100 { + t.Fatalf("pagesize 没有绑定成功: %d", yeeke.PageSize) + } + if yeeke.MaxPages != 10000 { + t.Fatalf("maxpages 没有绑定成功: %d", yeeke.MaxPages) + } + if yeeke.Retry != 2 { + t.Fatalf("retry 没有绑定成功: %d", yeeke.Retry) + } + if yeeke.OcrURL == "" { + t.Fatal("ocrurl 没有绑定成功") + } + if yeeke.OcrMaxAttempts != 5 { + t.Fatalf("ocrmaxattempts 没有绑定成功: %d", yeeke.OcrMaxAttempts) } } diff --git a/server/config/extend_test.go b/server/config/extend_test.go index e3c6bb1..9dde599 100644 --- a/server/config/extend_test.go +++ b/server/config/extend_test.go @@ -74,6 +74,61 @@ func TestApplyEnvironmentLoadsSYBCredentials(t *testing.T) { } } +// 同上,yeeke 凭据也只能来自环境变量 (#336)。 +func TestApplyEnvironmentLoadsYeekeCredentials(t *testing.T) { + original := ExtConfig.Yeeke + t.Cleanup(func() { ExtConfig.Yeeke = original }) + + t.Setenv("GOAUTO_YEEKE_USERNAME", " operator ") + t.Setenv("GOAUTO_YEEKE_PASSWORD", " se cret ") + ApplyEnvironment() + + if ExtConfig.Yeeke.Username != "operator" { + t.Fatalf("账号应去掉首尾空白: %q", ExtConfig.Yeeke.Username) + } + if ExtConfig.Yeeke.Password != " se cret " { + t.Fatalf("密码不应被修改: %q", ExtConfig.Yeeke.Password) + } +} + +func TestYeekeResolvedFillsBlanksButNeverInventsCredentials(t *testing.T) { + resolved := Yeeke{}.Resolved() + + if resolved.BaseURL != DefaultYeekeBaseURL { + t.Fatalf("BaseURL 默认值不对: %q", resolved.BaseURL) + } + if resolved.PageSize != DefaultYeekePageSize || resolved.MaxPages != DefaultYeekeMaxPages { + t.Fatalf("分页默认值不对: %+v", resolved) + } + if resolved.OcrMaxAttempts != DefaultYeekeOcrMaxAttempts { + t.Fatalf("OCR 重试次数默认值不对: %d", resolved.OcrMaxAttempts) + } + if resolved.Username != "" || resolved.Password != "" { + t.Fatal("Resolved 不得给账号密码编造默认值") + } + if resolved.OcrURL != "" { + t.Fatalf("空 OcrURL 不应被填充: %q", resolved.OcrURL) + } +} + +func TestYeekeHasCredentialsRequiresBothFields(t *testing.T) { + for _, c := range []struct { + name string + yeeke Yeeke + want bool + }{ + {"都有", Yeeke{Username: "a", Password: "b"}, true}, + {"缺密码", Yeeke{Username: "a"}, false}, + {"缺账号", Yeeke{Password: "b"}, false}, + {"账号只有空白", Yeeke{Username: " ", Password: "b"}, false}, + {"都没有", Yeeke{}, false}, + } { + if got := c.yeeke.HasCredentials(); got != c.want { + t.Fatalf("%s: 期望 %v,实际 %v", c.name, c.want, got) + } + } +} + func TestSYBResolvedFillsBlanksButNeverInventsCredentials(t *testing.T) { resolved := SYB{}.Resolved() diff --git a/server/config/local.go b/server/config/local.go index 40d4b94..77fafd5 100644 --- a/server/config/local.go +++ b/server/config/local.go @@ -65,6 +65,7 @@ type localFile struct { Database map[string]any `yaml:"database"` Ports map[string]any `yaml:"ports"` SYB map[string]any `yaml:"syb"` + Yeeke map[string]any `yaml:"yeeke"` } // ApplyLocalConfig loads config.yaml, if one is present, over the values @@ -96,6 +97,7 @@ func ApplyLocalConfig() { applyLocalDatabase(file.Database) applyLocalPorts(file.Ports) ApplyLocalSYB(file.SYB) + ApplyLocalYeeke(file.Yeeke) logInfo("本地配置:已加载 %s", path) } @@ -115,6 +117,17 @@ func LogEffectiveConfig() { } logInfo("顺云宝配置:凭据=%v(来源:%s)base_url=%s 验证码识别=%v", syb.HasCredentials(), source, syb.BaseURL, syb.OcrURL != "") + + yeeke := ExtConfig.Yeeke.Resolved() + yeekeSource := "未配置" + switch { + case strings.TrimSpace(os.Getenv("GOAUTO_YEEKE_USERNAME")) != "": + yeekeSource = "环境变量" + case yeeke.HasCredentials(): + yeekeSource = LocalConfigName + } + logInfo("yeeke 配置:凭据=%v(来源:%s)base_url=%s 验证码识别=%v", + yeeke.HasCredentials(), yeekeSource, yeeke.BaseURL, yeeke.OcrURL != "") } func applyLocalDatabase(database map[string]any) { @@ -177,6 +190,35 @@ func ApplyLocalSYB(syb map[string]any) { } } +// ApplyLocalYeeke folds a config.yaml `yeeke:` section into ExtConfig, mirroring +// ApplyLocalSYB above (#336). +func ApplyLocalYeeke(yeeke map[string]any) { + if username := strings.TrimSpace(scalar(yeeke, "username")); username != "" { + ExtConfig.Yeeke.Username = username + } + if password := scalar(yeeke, "password"); password != "" { + ExtConfig.Yeeke.Password = password + } + if baseURL := strings.TrimSpace(scalar(yeeke, "base_url")); baseURL != "" { + ExtConfig.Yeeke.BaseURL = baseURL + } + if ocrURL := strings.TrimSpace(scalar(yeeke, "ocr_url")); ocrURL != "" { + ExtConfig.Yeeke.OcrURL = ocrURL + } + if pageSize, err := strconv.Atoi(scalar(yeeke, "page_size")); err == nil && pageSize > 0 { + ExtConfig.Yeeke.PageSize = pageSize + } + if maxPages, err := strconv.Atoi(scalar(yeeke, "max_pages")); err == nil && maxPages > 0 { + ExtConfig.Yeeke.MaxPages = maxPages + } + if retry, err := strconv.Atoi(scalar(yeeke, "retry")); err == nil && retry >= 0 { + ExtConfig.Yeeke.Retry = retry + } + if attempts, err := strconv.Atoi(scalar(yeeke, "ocr_max_attempts")); err == nil && attempts > 0 { + ExtConfig.Yeeke.OcrMaxAttempts = attempts + } +} + // logInfo writes an informational startup line to stdout. // // `[必须]` Not stderr. The development launcher pipes the server through diff --git a/server/config/settings.yml b/server/config/settings.yml index 9c0a32e..70de3b6 100644 --- a/server/config/settings.yml +++ b/server/config/settings.yml @@ -61,6 +61,17 @@ settings: # `[必须]` 验证码图片会被发送到这个地址;换成别人运营的服务前要重新评估。 ocrurl: https://ocr.ilapage.cn/ocr ocrmaxattempts: 5 + # yeeke(mmt.yeeke.com)退货包裹只读对接,见 #336。 + # `[必须]` 账号密码不在这里,走环境变量 GOAUTO_YEEKE_USERNAME / GOAUTO_YEEKE_PASSWORD, + # 以免凭据进 Git。 + yeeke: + baseurl: https://mmt.yeeke.com + pagesize: 100 + maxpages: 10000 + retry: 2 + # 验证码识别服务,与 SYB 共用(#336 已批准)。 + ocrurl: https://ocr.ilapage.cn/ocr + ocrmaxattempts: 5 cache: # redis: # addr: 127.0.0.1:6379