Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
02bddbf304 | ||
|
|
3e82ad6570 | ||
|
|
b015719947 |
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 ""
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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},
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -44,6 +44,10 @@ func MigratedModels() []any {
|
||||
&models.SYBShop{},
|
||||
&models.SYBProductFilter{},
|
||||
&models.SYBSyncRun{},
|
||||
&models.YeekeSession{},
|
||||
&models.YeekeReturnPackage{},
|
||||
&models.YeekeReturnItem{},
|
||||
&models.YeekeSyncRun{},
|
||||
&models.SYBInnerCodeRecord{},
|
||||
&models.SYBInnerCodeItem{},
|
||||
&models.SYBInnerCodeApplyBatch{},
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
package models
|
||||
|
||||
import "time"
|
||||
|
||||
// YeekeSession stores only the opaque session material; credentials are kept
|
||||
// outside the application database and supplied by the administrator at run time.
|
||||
type YeekeSession struct {
|
||||
ID uint64 `gorm:"primaryKey;autoIncrement"`
|
||||
Username string `gorm:"size:128;not null;uniqueIndex:ux_yeeke_session_username"`
|
||||
Token string `json:"-" gorm:"type:text;not null"`
|
||||
CookiesJSON string `json:"-" gorm:"type:text;not null"`
|
||||
UserID string `gorm:"size:128;not null;default:''"`
|
||||
ExpiresAt time.Time `gorm:"not null;index"`
|
||||
CreatedAt time.Time
|
||||
UpdatedAt time.Time
|
||||
}
|
||||
|
||||
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:''"`
|
||||
// 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" }
|
||||
|
||||
type YeekeReturnItem struct {
|
||||
ID uint64 `gorm:"primaryKey;autoIncrement"`
|
||||
PackageID uint64 `gorm:"not null;uniqueIndex:ux_yeeke_return_item_key,priority:1;index"`
|
||||
ExternalKey string `gorm:"size:512;not null;uniqueIndex:ux_yeeke_return_item_key,priority:2"`
|
||||
ItemID string `gorm:"size:128;not null;index"`
|
||||
VariationID string `gorm:"size:128;not null;default:''"`
|
||||
ItemName string `gorm:"size:500;not null;default:''"`
|
||||
VariationName string `gorm:"size:500;not null;default:''"`
|
||||
Image string `gorm:"type:text;not null"`
|
||||
Quantity int64 `gorm:"not null;default:0"`
|
||||
LastSyncedAt time.Time `gorm:"not null;index"`
|
||||
SyncStatus string `gorm:"size:32;not null;default:'ok'"`
|
||||
CreatedAt time.Time
|
||||
UpdatedAt time.Time
|
||||
}
|
||||
|
||||
func (YeekeReturnItem) TableName() string { return "yeeke_return_item" }
|
||||
|
||||
type YeekeSyncRun struct {
|
||||
ID uint64 `gorm:"primaryKey;autoIncrement"`
|
||||
Status string `gorm:"size:32;not null;index"`
|
||||
Trigger string `gorm:"size:32;not null;index"`
|
||||
TotalPages int `gorm:"not null;default:0"`
|
||||
ReadCount int `gorm:"not null;default:0"`
|
||||
CreatedCount int `gorm:"not null;default:0"`
|
||||
UpdatedCount int `gorm:"not null;default:0"`
|
||||
SkippedCount int `gorm:"not null;default:0"`
|
||||
FailedCount int `gorm:"not null;default:0"`
|
||||
ErrorMessage string `gorm:"size:1000;not null;default:''"`
|
||||
StartedAt time.Time `gorm:"not null"`
|
||||
FinishedAt *time.Time
|
||||
LastSuccessAt *time.Time
|
||||
ActiveSlot *uint8 `gorm:"uniqueIndex:ux_yeeke_sync_run_active_slot"`
|
||||
LeaseOwner string `gorm:"size:128;not null;default:''"`
|
||||
LeaseExpiresAt *time.Time
|
||||
CreatedAt time.Time
|
||||
UpdatedAt time.Time
|
||||
}
|
||||
|
||||
func (YeekeSyncRun) TableName() string { return "yeeke_sync_run" }
|
||||
@@ -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}})
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -0,0 +1,301 @@
|
||||
package yeeke
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/app/goauto/yeekeclient"
|
||||
"gorm.io/gorm"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
PageSize, MaxPages, Retry int
|
||||
Lease time.Duration
|
||||
}
|
||||
|
||||
func (c Config) norm() Config {
|
||||
if c.PageSize <= 0 || c.PageSize > 500 {
|
||||
c.PageSize = 100
|
||||
}
|
||||
if c.MaxPages <= 0 || c.MaxPages > 10000 {
|
||||
c.MaxPages = 10000
|
||||
}
|
||||
if c.Retry < 0 || c.Retry > 5 {
|
||||
c.Retry = 2
|
||||
}
|
||||
if c.Lease <= 0 {
|
||||
c.Lease = 30 * time.Minute
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
type Report struct {
|
||||
RunID uint64
|
||||
TotalPages, Read, Created, Updated, Skipped, Failed int
|
||||
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() {
|
||||
return nil
|
||||
}
|
||||
x := t.Time
|
||||
return &x
|
||||
}
|
||||
func packageKey(p yeekeclient.ReturnPackage) string {
|
||||
if x := external(p.ID); x != "<nil>" && x != "" {
|
||||
return x
|
||||
}
|
||||
return p.Ordersn + "/" + p.TrackingNo + "/" + external(p.ShopID) + "/" + p.CreateTime.String()
|
||||
}
|
||||
|
||||
// itemKey builds the stable per-item identity used to upsert
|
||||
// models.YeekeReturnItem without duplicating rows across syncs. It prefers
|
||||
// the yeeke-issued identifiers (i.ID, then i.ItemID/i.VariationID) because
|
||||
// those stay the same regardless of the order the API returns items in
|
||||
// within a package; the positional index n is only used as a last resort
|
||||
// when none of those identifiers are present, since in that case the index
|
||||
// is the sole thing distinguishing items of the same package (#336).
|
||||
func itemKey(p yeekeclient.ReturnPackage, i yeekeclient.ReturnItem, n int) string {
|
||||
id, itemID, variationID := external(i.ID), external(i.ItemID), external(i.VariationID)
|
||||
if isEmptyExternal(id) && isEmptyExternal(itemID) && isEmptyExternal(variationID) {
|
||||
return packageKey(p) + "/" + id + "/" + itemID + "/" + variationID + "/" + strconv.Itoa(n)
|
||||
}
|
||||
return packageKey(p) + "/" + id + "/" + itemID + "/" + variationID
|
||||
}
|
||||
|
||||
// isEmptyExternal reports whether external() produced a value that carries
|
||||
// no real identity: either the field was unset (formatted as "<nil>" by
|
||||
// fmt.Sprint on a nil/zero value) or it was an explicit empty string.
|
||||
func isEmptyExternal(v string) bool { return v == "" || v == "<nil>" }
|
||||
|
||||
// takeoverStaleLease reclaims a run whose lease has expired, e.g. because the
|
||||
// process crashed or was restarted mid-sync. It matches the
|
||||
// lease-with-expiry-takeover idiom used by
|
||||
// app/goauto/purchase/order_writeback_worker.go: a single conditional UPDATE
|
||||
// guarded by "status = running AND lease_expires_at <= now" flips the stale
|
||||
// row to a terminal status and frees active_slot in one statement, so it is
|
||||
// atomic without a separate row lock. The stale row is never deleted — it is
|
||||
// left in place with status "failed" and an error_message explaining why, so
|
||||
// history stays auditable. If two callers race this same UPDATE, only the
|
||||
// first to reach the database actually changes any row; the second's WHERE
|
||||
// clause no longer matches (status is no longer "running") and it affects
|
||||
// zero rows, which is a harmless no-op. Whichever caller then wins the
|
||||
// subsequent Create (see acquire) is arbitrated by the ux_yeeke_sync_run_active_slot
|
||||
// unique index, exactly as it already is for two brand-new concurrent runs.
|
||||
func (s *Service) takeoverStaleLease(ctx context.Context) error {
|
||||
now := time.Now().UTC()
|
||||
return s.db.WithContext(ctx).Model(&models.YeekeSyncRun{}).
|
||||
Where("status = ? AND active_slot = ? AND lease_expires_at IS NOT NULL AND lease_expires_at <= ?", "running", 1, now).
|
||||
Updates(map[string]any{
|
||||
"status": "failed",
|
||||
"active_slot": nil,
|
||||
"lease_owner": "",
|
||||
"error_message": "lease expired: run interrupted, likely a process restart mid-sync (stale lease takeover)",
|
||||
"lease_expires_at": nil,
|
||||
"finished_at": now,
|
||||
}).Error
|
||||
}
|
||||
|
||||
func (s *Service) acquire(ctx context.Context, trigger string) (*models.YeekeSyncRun, error) {
|
||||
if e := s.takeoverStaleLease(ctx); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
owner := fmt.Sprintf("%d", now.UnixNano())
|
||||
slot := uint8(1)
|
||||
exp := now.Add(s.cfg.Lease)
|
||||
r := &models.YeekeSyncRun{Status: "running", Trigger: trigger, StartedAt: now, ActiveSlot: &slot, LeaseOwner: owner, LeaseExpiresAt: &exp}
|
||||
if e := s.db.WithContext(ctx).Create(r).Error; e != nil {
|
||||
return nil, e
|
||||
}
|
||||
return r, nil
|
||||
}
|
||||
|
||||
type Service struct {
|
||||
db *gorm.DB
|
||||
client *yeekeclient.Client
|
||||
cfg Config
|
||||
}
|
||||
|
||||
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()
|
||||
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 {
|
||||
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(e.Error())
|
||||
return rep, runErr
|
||||
}
|
||||
rep.TotalPages = page
|
||||
if len(p.Records) == 0 {
|
||||
break
|
||||
}
|
||||
finger := pageFingerprint(p)
|
||||
if seen[finger] {
|
||||
rep.Skipped += len(p.Records)
|
||||
break
|
||||
}
|
||||
seen[finger] = true
|
||||
for _, x := range p.Records {
|
||||
created, updated, err := s.upsert(ctx, x)
|
||||
if err != nil {
|
||||
rep.Failed++
|
||||
continue
|
||||
}
|
||||
rep.Read++
|
||||
if created {
|
||||
rep.Created++
|
||||
} else if updated {
|
||||
rep.Updated++
|
||||
} else {
|
||||
rep.Skipped++
|
||||
}
|
||||
}
|
||||
if len(p.Records) < s.cfg.PageSize {
|
||||
break
|
||||
}
|
||||
if p.Pages > 0 && page >= p.Pages {
|
||||
break
|
||||
}
|
||||
}
|
||||
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)
|
||||
return hex.EncodeToString(h[:])
|
||||
}
|
||||
func (s *Service) upsert(ctx context.Context, p yeekeclient.ReturnPackage) (bool, bool, error) {
|
||||
now := time.Now().UTC()
|
||||
key := packageKey(p)
|
||||
var row models.YeekeReturnPackage
|
||||
e := s.db.WithContext(ctx).Where("external_id = ?", key).First(&row).Error
|
||||
isNew := e == gorm.ErrRecordNotFound
|
||||
if e != nil && !isNew {
|
||||
return false, false, e
|
||||
}
|
||||
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 {
|
||||
return false, false, e
|
||||
}
|
||||
}
|
||||
if e = s.db.WithContext(ctx).Model(&row).Updates(vals).Error; e != nil {
|
||||
return false, false, e
|
||||
}
|
||||
for n, i := range p.Items {
|
||||
ik := itemKey(p, i, n)
|
||||
ir := models.YeekeReturnItem{}
|
||||
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.WithContext(ctx).Create(&models.YeekeReturnItem{PackageID: row.ID, ExternalKey: ik}).Error; e != nil {
|
||||
return false, false, e
|
||||
}
|
||||
}
|
||||
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
|
||||
}
|
||||
}
|
||||
return isNew, !isNew, nil
|
||||
}
|
||||
@@ -0,0 +1,730 @@
|
||||
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 }
|
||||
}
|
||||
|
||||
// packageWithItems builds a raw list-page record for one package carrying an
|
||||
// arbitrary, caller-ordered set of items, so tests can reorder items between
|
||||
// two syncs of the same package.
|
||||
func packageWithItems(pkgID string, items ...map[string]any) string {
|
||||
b, _ := json.Marshal(map[string]any{
|
||||
"id": pkgID, "ordersn": "o-" + pkgID, "trackingNo": "t-" + pkgID, "status": 1,
|
||||
"items": items,
|
||||
})
|
||||
return string(b)
|
||||
}
|
||||
|
||||
func item(id, itemID, variationID string) map[string]any {
|
||||
return map[string]any{"id": id, "itemId": itemID, "variationId": variationID, "itemName": "n", "variationName": "v", "variationQuantityPurchased": 1}
|
||||
}
|
||||
|
||||
// itemNoIDs builds an item carrying no yeeke-issued identifiers at all
|
||||
// (id/itemId/variationId all empty), the case itemKey's positional-index
|
||||
// fallback exists for.
|
||||
func itemNoIDs() map[string]any {
|
||||
return map[string]any{"id": "", "itemId": "", "variationId": "", "itemName": "n", "variationName": "v", "variationQuantityPurchased": 1}
|
||||
}
|
||||
|
||||
// --- Defect 1 (#336): stale lease takeover -------------------------------
|
||||
|
||||
// TestStaleLeaseIsTakenOverOnNextAcquire simulates a crash: a "running" row
|
||||
// is left behind with an active_slot and a lease that has already expired
|
||||
// (as if the process died mid-sync, long before the lease's normal
|
||||
// duration). The very next Sync call — scheduled or manual — must reclaim
|
||||
// the slot rather than being permanently blocked, and the abandoned row must
|
||||
// end up in a clear terminal state (not silently deleted) recording why.
|
||||
func TestStaleLeaseIsTakenOverOnNextAcquire(t *testing.T) {
|
||||
db := testDB(t)
|
||||
slot := uint8(1)
|
||||
past := time.Now().UTC().Add(-time.Hour)
|
||||
stale := models.YeekeSyncRun{
|
||||
Status: "running", Trigger: "scheduled", StartedAt: past.Add(-time.Minute),
|
||||
ActiveSlot: &slot, LeaseOwner: "dead-process", LeaseExpiresAt: &past,
|
||||
}
|
||||
if e := db.Create(&stale).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
|
||||
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.Fatalf("resume after a stale lease must succeed, got err=%v", err)
|
||||
}
|
||||
if rep.Status != "succeeded" {
|
||||
t.Fatalf("rep=%+v", rep)
|
||||
}
|
||||
|
||||
var reclaimed models.YeekeSyncRun
|
||||
if e := db.First(&reclaimed, stale.ID).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
if reclaimed.Status != "failed" {
|
||||
t.Fatalf("stale run status=%q, want a terminal status (not silently left running or deleted)", reclaimed.Status)
|
||||
}
|
||||
if reclaimed.ErrorMessage == "" {
|
||||
t.Fatal("stale run must record why it was taken over")
|
||||
}
|
||||
if reclaimed.ActiveSlot != nil {
|
||||
t.Fatal("stale run must release active_slot on takeover")
|
||||
}
|
||||
|
||||
// The new run's own row clears active_slot on completion just like any
|
||||
// other successful run (see run()'s defer), so what proves the takeover
|
||||
// happened is that a second, distinct run row now exists alongside the
|
||||
// reclaimed stale one.
|
||||
var totalRuns int64
|
||||
db.Model(&models.YeekeSyncRun{}).Count(&totalRuns)
|
||||
if totalRuns != 2 {
|
||||
t.Fatalf("expected the stale row plus exactly one new run after takeover, got %d run rows", totalRuns)
|
||||
}
|
||||
if rep.RunID == stale.ID {
|
||||
t.Fatal("the new run must not reuse the stale run's row")
|
||||
}
|
||||
}
|
||||
|
||||
// TestActiveSlotLeaseRejectsConcurrentRuns above must still pass unmodified:
|
||||
// a lease that has NOT expired must keep blocking a second run. This test
|
||||
// pins that same guarantee at the acquire() level directly.
|
||||
func TestValidLeaseIsNotTakenOver(t *testing.T) {
|
||||
db := testDB(t)
|
||||
slot := uint8(1)
|
||||
future := time.Now().UTC().Add(time.Hour)
|
||||
holding := models.YeekeSyncRun{
|
||||
Status: "running", Trigger: "manual", StartedAt: time.Now().UTC(),
|
||||
ActiveSlot: &slot, LeaseOwner: "still-alive", LeaseExpiresAt: &future,
|
||||
}
|
||||
if e := db.Create(&holding).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
c, _ := yeekeclient.New("http://unused.invalid")
|
||||
s := NewService(db, c, Config{PageSize: 10})
|
||||
if _, e := s.acquire(context.Background(), "scheduled"); e == nil {
|
||||
t.Fatal("a still-valid lease must not be taken over or bypassed")
|
||||
}
|
||||
var row models.YeekeSyncRun
|
||||
if e := db.First(&row, holding.ID).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
if row.Status != "running" || row.ActiveSlot == nil {
|
||||
t.Fatalf("holder must be untouched: %+v", row)
|
||||
}
|
||||
}
|
||||
|
||||
// TestConcurrentTakeoverExactlyOneWins races two acquire() calls against the
|
||||
// same stale, expired-lease row. Both attempt the takeover UPDATE and then a
|
||||
// Create; the takeover UPDATE is idempotent (the loser affects zero rows
|
||||
// since the row's status is no longer "running" by the time it runs), and
|
||||
// the ux_yeeke_sync_run_active_slot unique index arbitrates the Create race
|
||||
// the same way it already does for two brand-new concurrent runs. Exactly
|
||||
// one goroutine must come away holding the slot.
|
||||
func TestConcurrentTakeoverExactlyOneWins(t *testing.T) {
|
||||
db := testDB(t)
|
||||
slot := uint8(1)
|
||||
past := time.Now().UTC().Add(-time.Hour)
|
||||
stale := models.YeekeSyncRun{
|
||||
Status: "running", Trigger: "scheduled", StartedAt: past.Add(-time.Minute),
|
||||
ActiveSlot: &slot, LeaseOwner: "dead-process", LeaseExpiresAt: &past,
|
||||
}
|
||||
if e := db.Create(&stale).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
// SQLite only allows one writer at a time; serialize connections through
|
||||
// the Go pool (same pattern as app/goauto/purchase/order_backfill_test.go
|
||||
// and friends) so the race is decided by acquire()'s own logic rather
|
||||
// than by spurious "database is locked" errors.
|
||||
if sqlDB, e := db.DB(); e == nil {
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
}
|
||||
c, _ := yeekeclient.New("http://unused.invalid")
|
||||
s := NewService(db, c, Config{PageSize: 10})
|
||||
|
||||
const n = 8
|
||||
var wg sync.WaitGroup
|
||||
oks := make([]bool, n)
|
||||
for i := 0; i < n; i++ {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
_, e := s.acquire(context.Background(), "manual")
|
||||
oks[i] = e == nil
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
winners := 0
|
||||
for _, ok := range oks {
|
||||
if ok {
|
||||
winners++
|
||||
}
|
||||
}
|
||||
if winners != 1 {
|
||||
t.Fatalf("winners=%d, want exactly 1 (active_slot must arbitrate concurrent takeover attempts)", winners)
|
||||
}
|
||||
var holders int64
|
||||
db.Model(&models.YeekeSyncRun{}).Where("active_slot = ?", 1).Count(&holders)
|
||||
if holders != 1 {
|
||||
t.Fatalf("holders=%d, want exactly 1 row holding active_slot after the race", holders)
|
||||
}
|
||||
}
|
||||
|
||||
// --- Defect 2 (#336): itemKey must not depend on item order --------------
|
||||
|
||||
// TestItemKeyStableAcrossReorder syncs the same package twice with its two
|
||||
// items in reversed order the second time. Reordering must not create new
|
||||
// rows: each item's identity must key off its own IDs, not its position.
|
||||
func TestItemKeyStableAcrossReorder(t *testing.T) {
|
||||
db := testDB(t)
|
||||
var call int32
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
n := atomic.AddInt32(&call, 1)
|
||||
if n == 1 {
|
||||
fmt.Fprint(w, page([]string{packageWithItems("p1", item("i1", "item-a", "var-a"), item("i2", "item-b", "var-b"))}, 1, 1))
|
||||
return
|
||||
}
|
||||
if n == 2 {
|
||||
// Same package, items reordered.
|
||||
fmt.Fprint(w, page([]string{packageWithItems("p1", item("i2", "item-b", "var-b"), item("i1", "item-a", "var-a"))}, 1, 1))
|
||||
return
|
||||
}
|
||||
fmt.Fprint(w, page(nil, 0, 0))
|
||||
}))
|
||||
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)
|
||||
}
|
||||
if _, err := s.Sync(context.Background(), "manual"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
var pkg models.YeekeReturnPackage
|
||||
if e := db.Where("external_id = ?", "p1").First(&pkg).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
var n int64
|
||||
db.Model(&models.YeekeReturnItem{}).Where("package_id = ?", pkg.ID).Count(&n)
|
||||
if n != 2 {
|
||||
t.Fatalf("items=%d, want 2 (reordering the same items must not duplicate rows)", n)
|
||||
}
|
||||
}
|
||||
|
||||
// TestItemKeyIndexFallbackForItemsLackingAllIDs covers a package whose items
|
||||
// carry no yeeke-issued identifiers at all: the positional index is the only
|
||||
// thing that can distinguish them, so the fallback must still apply and keep
|
||||
// them as separate rows.
|
||||
func TestItemKeyIndexFallbackForItemsLackingAllIDs(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{packageWithItems("p1", itemNoIDs(), itemNoIDs())}, 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 pkg models.YeekeReturnPackage
|
||||
if e := db.Where("external_id = ?", "p1").First(&pkg).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
var n int64
|
||||
db.Model(&models.YeekeReturnItem{}).Where("package_id = ?", pkg.ID).Count(&n)
|
||||
if n != 2 {
|
||||
t.Fatalf("items=%d, want 2 (items lacking all IDs must still be distinguished by position)", n)
|
||||
}
|
||||
}
|
||||
|
||||
// TestItemKeyDistinctVariationsOfSameItemID pins existing behavior: two
|
||||
// items sharing the same itemID but different variationIDs are, and must
|
||||
// remain, two distinct rows.
|
||||
func TestItemKeyDistinctVariationsOfSameItemID(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{packageWithItems("p1", item("i1", "item-a", "var-1"), item("i2", "item-a", "var-2"))}, 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 pkg models.YeekeReturnPackage
|
||||
if e := db.Where("external_id = ?", "p1").First(&pkg).Error; e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
var n int64
|
||||
db.Model(&models.YeekeReturnItem{}).Where("package_id = ?", pkg.ID).Count(&n)
|
||||
if n != 2 {
|
||||
t.Fatalf("items=%d, want 2 (same itemID with different variationID must stay distinct)", n)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
package yeeke
|
||||
|
||||
import (
|
||||
"context"
|
||||
"go-admin/app/goauto/migrations"
|
||||
"go-admin/app/goauto/models"
|
||||
"go-admin/app/goauto/yeekeclient"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestSyncIsIdempotentAndKeepsVariationsSeparate(t *testing.T) {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.Write([]byte(`{"success":true,"result":{"records":[{"id":"p1","ordersn":"o","trackingNo":"t","status":1,"items":[{"id":"a","itemId":"i","variationId":"v1","itemName":"n","variationName":"red","variationQuantityPurchased":1},{"id":"b","itemId":"i","variationId":"v2","itemName":"n","variationName":"blue","variationQuantityPurchased":1}]}],"total":1,"pages":1}}`))
|
||||
}))
|
||||
defer server.Close()
|
||||
db, _ := gorm.Open(sqlite.Open("file:yeeke-sync?mode=memory&cache=shared"), &gorm.Config{})
|
||||
if e := migrations.Migrate(db); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
c, _ := yeekeclient.New(server.URL)
|
||||
s := NewService(db, c, Config{PageSize: 10})
|
||||
if _, e := s.Sync(context.Background(), "manual"); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
if _, e := s.Sync(context.Background(), "manual"); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
var n int64
|
||||
db.Model(&models.YeekeReturnPackage{}).Count(&n)
|
||||
if n != 1 {
|
||||
t.Fatalf("packages=%d", n)
|
||||
}
|
||||
db.Model(&models.YeekeReturnItem{}).Count(&n)
|
||||
if n != 2 {
|
||||
t.Fatalf("items=%d", n)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,294 @@
|
||||
package yeekeclient
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/cookiejar"
|
||||
"net/url"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
var ErrSessionInvalid = errors.New("yeeke session invalid")
|
||||
var ErrNoSession = errors.New("yeeke session unavailable")
|
||||
|
||||
type Session struct {
|
||||
Username, Token, CookiesJSON, UserID string
|
||||
ExpiresAt time.Time
|
||||
}
|
||||
type Captcha struct {
|
||||
Image []byte
|
||||
CheckKey string
|
||||
ContentType string
|
||||
}
|
||||
type Client struct {
|
||||
baseURL string
|
||||
http *http.Client
|
||||
jar *cookiejar.Jar
|
||||
token string
|
||||
retry int
|
||||
}
|
||||
|
||||
func New(baseURL string) (*Client, error) {
|
||||
baseURL = strings.TrimRight(strings.TrimSpace(baseURL), "/")
|
||||
if baseURL == "" {
|
||||
return nil, fmt.Errorf("yeeke base_url required")
|
||||
}
|
||||
j, e := cookiejar.New(nil)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
return &Client{baseURL: baseURL, jar: j, http: &http.Client{Jar: j, Timeout: 60 * time.Second}}, nil
|
||||
}
|
||||
func (c *Client) SetToken(t string) { c.token = t }
|
||||
func (c *Client) Token() string { return c.token }
|
||||
|
||||
type cookieDTO struct{ Name, Value, Path string }
|
||||
|
||||
func (c *Client) ExportCookiesJSON() (string, error) {
|
||||
u, e := url.Parse(c.baseURL)
|
||||
if e != nil {
|
||||
return "", e
|
||||
}
|
||||
a := []cookieDTO{}
|
||||
for _, x := range c.jar.Cookies(u) {
|
||||
a = append(a, cookieDTO{x.Name, x.Value, x.Path})
|
||||
}
|
||||
b, e := json.Marshal(a)
|
||||
return string(b), e
|
||||
}
|
||||
func (c *Client) ImportCookiesJSON(s string) error {
|
||||
var a []cookieDTO
|
||||
if e := json.Unmarshal([]byte(s), &a); e != nil {
|
||||
return e
|
||||
}
|
||||
u, e := url.Parse(c.baseURL)
|
||||
if e != nil {
|
||||
return e
|
||||
}
|
||||
cs := []*http.Cookie{}
|
||||
for _, x := range a {
|
||||
if x.Name != "" {
|
||||
p := x.Path
|
||||
if p == "" {
|
||||
p = "/"
|
||||
}
|
||||
cs = append(cs, &http.Cookie{Name: x.Name, Value: x.Value, Path: p})
|
||||
}
|
||||
}
|
||||
c.jar.SetCookies(u, cs)
|
||||
return nil
|
||||
}
|
||||
func (c *Client) do(ctx context.Context, method, path string, body any, query url.Values) (json.RawMessage, error) {
|
||||
b := io.Reader(nil)
|
||||
if body != nil {
|
||||
x, e := json.Marshal(body)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
b = bytes.NewReader(x)
|
||||
}
|
||||
u := c.baseURL + path
|
||||
if len(query) > 0 {
|
||||
u += "?" + query.Encode()
|
||||
}
|
||||
req, e := http.NewRequestWithContext(ctx, method, u, b)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
req.Header.Set("Accept", "application/json")
|
||||
if body != nil {
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
}
|
||||
if c.token != "" {
|
||||
q := req.URL.Query()
|
||||
q.Set("token", c.token)
|
||||
req.URL.RawQuery = q.Encode()
|
||||
}
|
||||
resp, e := c.http.Do(req)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
raw, e := io.ReadAll(resp.Body)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
if resp.StatusCode == 401 || resp.StatusCode == 403 {
|
||||
return nil, ErrSessionInvalid
|
||||
}
|
||||
if resp.StatusCode >= 500 {
|
||||
return nil, fmt.Errorf("yeeke http %d", resp.StatusCode)
|
||||
}
|
||||
var env struct {
|
||||
Success bool `json:"success"`
|
||||
Code int `json:"code"`
|
||||
Message string `json:"message"`
|
||||
Result json.RawMessage `json:"result"`
|
||||
}
|
||||
if e = json.Unmarshal(raw, &env); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
if !env.Success {
|
||||
if env.Code == 401 || strings.Contains(env.Message, "登录") || strings.Contains(env.Message, "token") {
|
||||
return nil, ErrSessionInvalid
|
||||
}
|
||||
return nil, fmt.Errorf("yeeke request failed code=%d", env.Code)
|
||||
}
|
||||
return env.Result, nil
|
||||
}
|
||||
func (c *Client) FetchCaptcha(ctx context.Context) (*Captcha, error) {
|
||||
raw, e := c.do(ctx, http.MethodGet, "/agent-foreign/sys/randomImage", nil, nil)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
var p struct {
|
||||
Image string `json:"image"`
|
||||
CheckKey string `json:"checkKey"`
|
||||
}
|
||||
if e = json.Unmarshal(raw, &p); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
s := p.Image
|
||||
if i := strings.Index(s, ","); i >= 0 {
|
||||
s = s[i+1:]
|
||||
}
|
||||
img, e := base64.StdEncoding.DecodeString(s)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
return &Captcha{Image: img, CheckKey: p.CheckKey, ContentType: "image/jpeg"}, nil
|
||||
}
|
||||
|
||||
type LoginResult struct {
|
||||
Token, UserID, Username string
|
||||
ExpiresAt time.Time
|
||||
}
|
||||
|
||||
type OCR interface {
|
||||
Recognize(context.Context, []byte) (string, error)
|
||||
}
|
||||
|
||||
// LoginWithOCR keeps captcha bytes in memory and never includes credentials or
|
||||
// recognized text in returned errors. A fresh image is fetched for every try.
|
||||
func (c *Client) LoginWithOCR(ctx context.Context, ocr OCR, username, password string, maxAttempts int) (*LoginResult, error) {
|
||||
if ocr == nil {
|
||||
return nil, fmt.Errorf("yeeke OCR unavailable")
|
||||
}
|
||||
if maxAttempts <= 0 || maxAttempts > 5 {
|
||||
maxAttempts = 3
|
||||
}
|
||||
for i := 0; i < maxAttempts; i++ {
|
||||
cap, e := c.FetchCaptcha(ctx)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
code, e := ocr.Recognize(ctx, cap.Image)
|
||||
if e != nil {
|
||||
return nil, fmt.Errorf("yeeke OCR unavailable")
|
||||
}
|
||||
code = strings.TrimSpace(code)
|
||||
if code == "" {
|
||||
continue
|
||||
}
|
||||
if out, e := c.Login(ctx, username, password, code, cap.CheckKey); e == nil {
|
||||
return out, nil
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("yeeke login failed after limited captcha attempts")
|
||||
}
|
||||
|
||||
func (c *Client) Login(ctx context.Context, username, password, captcha, checkKey string) (*LoginResult, error) {
|
||||
if username == "" || password == "" || captcha == "" || checkKey == "" {
|
||||
return nil, fmt.Errorf("login fields required")
|
||||
}
|
||||
raw, e := c.do(ctx, http.MethodPost, "/agent-foreign/sys/login", map[string]string{"username": username, "password": password, "captcha": captcha, "checkKey": checkKey, "agentCode": "mmt"}, nil)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
var p struct {
|
||||
Token string `json:"token"`
|
||||
UserInfo struct {
|
||||
ID any `json:"id"`
|
||||
Username string `json:"username"`
|
||||
} `json:"userInfo"`
|
||||
}
|
||||
if e = json.Unmarshal(raw, &p); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
if p.Token == "" {
|
||||
return nil, fmt.Errorf("yeeke login response missing token")
|
||||
}
|
||||
c.token = p.Token
|
||||
return &LoginResult{Token: p.Token, UserID: fmt.Sprint(p.UserInfo.ID), Username: p.UserInfo.Username, ExpiresAt: time.Now().UTC().Add(24 * time.Hour)}, nil
|
||||
}
|
||||
func (c *Client) CheckSession(ctx context.Context) error {
|
||||
_, e := c.do(ctx, http.MethodGet, "/agent-foreign/sys/userInfo", nil, nil)
|
||||
return e
|
||||
}
|
||||
|
||||
type ReturnPage struct {
|
||||
Records []ReturnPackage `json:"records"`
|
||||
Total int `json:"total"`
|
||||
Pages int `json:"pages"`
|
||||
}
|
||||
type Timestamp struct{ time.Time }
|
||||
|
||||
func (t *Timestamp) UnmarshalJSON(b []byte) error {
|
||||
var s string
|
||||
if json.Unmarshal(b, &s) != nil || s == "" {
|
||||
return nil
|
||||
}
|
||||
for _, f := range []string{time.RFC3339, "2006-01-02 15:04:05", "2006-01-02"} {
|
||||
if x, e := time.ParseInLocation(f, s, time.UTC); e == nil {
|
||||
t.Time = x
|
||||
return nil
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type ReturnPackage struct {
|
||||
ID any `json:"id"`
|
||||
Ordersn string `json:"ordersn"`
|
||||
TrackingNo string `json:"trackingNo"`
|
||||
ShopID any `json:"shopId"`
|
||||
ShopName string `json:"shopName"`
|
||||
WareCode string `json:"wareCode"`
|
||||
WareHouse string `json:"wareHouse"`
|
||||
WareName string `json:"wareName"`
|
||||
Status any `json:"status"`
|
||||
ClaimTime *Timestamp `json:"claimTime"`
|
||||
CreateTime *Timestamp `json:"createTime"`
|
||||
UpdateTime *Timestamp `json:"updateTime"`
|
||||
DestroyDeadLine *Timestamp `json:"destroyDeadLine"`
|
||||
Items []ReturnItem `json:"items"`
|
||||
}
|
||||
type ReturnItem struct {
|
||||
ID any `json:"id"`
|
||||
ItemID any `json:"itemId"`
|
||||
VariationID any `json:"variationId"`
|
||||
ItemName string `json:"itemName"`
|
||||
VariationName string `json:"variationName"`
|
||||
Image string `json:"image"`
|
||||
Quantity int64 `json:"variationQuantityPurchased"`
|
||||
}
|
||||
|
||||
func (c *Client) List(ctx context.Context, pageNo, pageSize int) (ReturnPage, error) {
|
||||
body := map[string]any{"pageNo": pageNo, "pageSize": pageSize, "claimFlag": 1, "status": 1, "relationFlag": 1, "orderBy": "createTime", "order": "desc"}
|
||||
raw, e := c.do(ctx, http.MethodPost, "/agent-foreign/packageClaimRec/relation/list", body, nil)
|
||||
if e != nil {
|
||||
return ReturnPage{}, e
|
||||
}
|
||||
var p ReturnPage
|
||||
if e = json.Unmarshal(raw, &p); e != nil {
|
||||
return p, e
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
package yeekeclient
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestCaptchaLoginAndReadOnlyList(t *testing.T) {
|
||||
s := 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": true, "result": map[string]any{"token": "opaque", "userInfo": map[string]any{"id": "u"}}})
|
||||
case "/agent-foreign/packageClaimRec/relation/list":
|
||||
if r.URL.Query().Get("token") != "opaque" {
|
||||
t.Errorf("token missing")
|
||||
}
|
||||
json.NewEncoder(w).Encode(map[string]any{"success": true, "result": map[string]any{"records": []any{}, "total": 0}})
|
||||
default:
|
||||
http.NotFound(w, r)
|
||||
}
|
||||
}))
|
||||
defer s.Close()
|
||||
c, _ := New(s.URL)
|
||||
cap, e := c.FetchCaptcha(context.Background())
|
||||
if e != nil || string(cap.Image) != "Hi" || cap.CheckKey != "k" {
|
||||
t.Fatalf("captcha=%+v err=%v", cap, e)
|
||||
}
|
||||
if _, e = c.Login(context.Background(), "u", "p", "1234", "k"); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
if _, e = c.List(context.Background(), 1, 10); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
package yeekeclient
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Credentials struct{ Username, Password string }
|
||||
|
||||
// Connect restores and validates a cached session. Only an explicit invalid
|
||||
// response deletes it; timeouts and 5xx preserve the usable cache.
|
||||
func Connect(ctx context.Context, store *SessionStore, creds Credentials, baseURL string, ocr OCR, maxLogin int) (*Client, error) {
|
||||
c, e := New(baseURL)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
s, e := store.Load(ctx, creds.Username, time.Now().UTC())
|
||||
if e == nil {
|
||||
if e = c.ImportCookiesJSON(s.CookiesJSON); e == nil {
|
||||
c.SetToken(s.Token)
|
||||
if e = c.CheckSession(ctx); e == nil {
|
||||
return c, nil
|
||||
} else if errors.Is(e, ErrSessionInvalid) {
|
||||
_ = store.Delete(ctx, creds.Username)
|
||||
} else {
|
||||
return nil, e
|
||||
}
|
||||
}
|
||||
} else if !errors.Is(e, ErrNoSession) {
|
||||
return nil, e
|
||||
}
|
||||
r, e := c.LoginWithOCR(ctx, ocr, creds.Username, creds.Password, maxLogin)
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
cookies, _ := c.ExportCookiesJSON()
|
||||
if e = store.Save(ctx, Session{Username: creds.Username, Token: r.Token, UserID: r.UserID, CookiesJSON: cookies, ExpiresAt: r.ExpiresAt}); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
return c, nil
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package yeekeclient
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"go-admin/app/goauto/models"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
type SessionStore struct{ db *gorm.DB }
|
||||
|
||||
func NewSessionStore(db *gorm.DB) *SessionStore { return &SessionStore{db: db} }
|
||||
func (s *SessionStore) Save(ctx context.Context, x Session) error {
|
||||
if strings.TrimSpace(x.Username) == "" || x.Token == "" || x.ExpiresAt.IsZero() {
|
||||
return errors.New("invalid yeeke session")
|
||||
}
|
||||
return s.db.WithContext(ctx).Clauses(clause.OnConflict{Columns: []clause.Column{{Name: "username"}}, DoUpdates: clause.AssignmentColumns([]string{"token", "cookies_json", "user_id", "expires_at", "updated_at"})}).Create(&models.YeekeSession{Username: strings.TrimSpace(x.Username), Token: x.Token, CookiesJSON: x.CookiesJSON, UserID: x.UserID, ExpiresAt: x.ExpiresAt.UTC()}).Error
|
||||
}
|
||||
func (s *SessionStore) Load(ctx context.Context, user string, now time.Time) (Session, error) {
|
||||
var r models.YeekeSession
|
||||
if e := s.db.WithContext(ctx).Where("username = ?", strings.TrimSpace(user)).First(&r).Error; e != nil {
|
||||
if errors.Is(e, gorm.ErrRecordNotFound) {
|
||||
return Session{}, ErrNoSession
|
||||
}
|
||||
return Session{}, e
|
||||
}
|
||||
if !now.UTC().Before(r.ExpiresAt) {
|
||||
return Session{}, ErrNoSession
|
||||
}
|
||||
return Session{Username: r.Username, Token: r.Token, CookiesJSON: r.CookiesJSON, UserID: r.UserID, ExpiresAt: r.ExpiresAt}, nil
|
||||
}
|
||||
func (s *SessionStore) Delete(ctx context.Context, user string) error {
|
||||
return s.db.WithContext(ctx).Where("username = ?", strings.TrimSpace(user)).Delete(&models.YeekeSession{}).Error
|
||||
}
|
||||
@@ -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{},
|
||||
// ...
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
package version_local
|
||||
|
||||
import (
|
||||
"go-admin/app/goauto/migrations"
|
||||
"go-admin/cmd/migrate/migration"
|
||||
common "go-admin/common/models"
|
||||
"gorm.io/gorm"
|
||||
"runtime"
|
||||
)
|
||||
|
||||
func init() {
|
||||
_, f, _, _ := runtime.Caller(0)
|
||||
migration.Migrate.SetVersion(migration.GetFilename(f), migrateYeekeReturnSync)
|
||||
}
|
||||
func migrateYeekeReturnSync(db *gorm.DB, version string) error {
|
||||
return db.Transaction(func(tx *gorm.DB) error {
|
||||
if err := migrations.Migrate(tx); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(&common.Migration{Version: version}).Error
|
||||
})
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
+69
-2
@@ -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 == "" {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user