diff --git a/admin/handler/web/others.go b/admin/handler/web/others.go index c653d45..d805ac9 100644 --- a/admin/handler/web/others.go +++ b/admin/handler/web/others.go @@ -97,6 +97,11 @@ func (h *Handler) renderSybListWithLoginReason(c *gin.Context, keyword, pageRaw, fail(c, http.StatusInternalServerError, "读取可分配客户端失败,数据没有被改动。") return } + enabledShopCount, err := service.CountEnabledSybAllowedShops(h.db) + if err != nil { + fail(c, http.StatusInternalServerError, "读取顺运宝同步店铺数量失败,数据没有被改动。刷新后重试。") + return + } syncStatus := service.GetSybSyncStatus() status := msg @@ -217,10 +222,11 @@ func (h *Handler) renderSybListWithLoginReason(c *gin.Context, keyword, pageRaw, keyword, shop, stage, dateFrom, dateTo), "HistoryFetchPagination": sybHistoryPartialPagination(history.Page, history.TotalPages, keyword, shop, stage, dateFrom, dateTo), - "Pagination": service.NewPaginationView(result.Page, result.TotalPages, values.Encode()), - "DetailURL": "/syb/detail?" + detailValues.Encode(), - "AssignableClients": purchaseClients.Rows, - "PurchaseClientCount": purchaseClients.SelectableCount, + "Pagination": service.NewPaginationView(result.Page, result.TotalPages, values.Encode()), + "DetailURL": "/syb/detail?" + detailValues.Encode(), + "AssignableClients": purchaseClients.Rows, + "PurchaseClientCount": purchaseClients.SelectableCount, + "EnabledSyncShopCount": enabledShopCount, })) } @@ -363,6 +369,10 @@ func (h *Handler) SybSync(c *gin.Context) { h.sybRedirectRange(c, err.Error()) return } + if err := service.EnsureEnabledSybAllowedShops(h.db); err != nil { + h.sybRedirect(c, "同步没有启动:"+err.Error()) + return + } cfg, err := config.Load() if err != nil { @@ -481,6 +491,10 @@ func (h *Handler) SybLoginAndSync(c *gin.Context) { h.sybRedirectRange(c, optionErr.Error()) return } + if err := service.EnsureEnabledSybAllowedShops(h.db); err != nil { + h.sybRedirect(c, "同步没有启动:"+err.Error()) + return + } cfg, err := config.Load() if err != nil { diff --git a/admin/handler/web/syb_shop.go b/admin/handler/web/syb_shop.go new file mode 100644 index 0000000..e737287 --- /dev/null +++ b/admin/handler/web/syb_shop.go @@ -0,0 +1,76 @@ +package web + +import ( + "errors" + "net/http" + "net/url" + "time" + + "github.com/gin-gonic/gin" + + "cmautobuy/admin/repository" + "cmautobuy/admin/service" +) + +func (h *Handler) SybAllowedShopList(c *gin.Context) { + rows, err := service.ListSybAllowedShops(h.db, currentUser(c)) + if err != nil { + fail(c, http.StatusInternalServerError, "读取顺运宝同步店铺失败,数据没有被改动。刷新后重试。") + return + } + enabled := 0 + for _, row := range rows { + if row.Enabled { + enabled++ + } + } + c.HTML(http.StatusOK, "syb/shops", page(c, "syb", "顺运宝同步店铺", gin.H{ + "Rows": rows, "EnabledCount": enabled, "Message": c.Query("msg"), "Error": c.Query("error"), + })) +} + +func (h *Handler) SybAllowedShopCreate(c *gin.Context) { + err := service.CreateSybAllowedShop(h.db, currentUser(c), c.PostForm("shop_name"), time.Now()) + if err != nil { + if service.IsValidationError(err) || errors.Is(err, repository.ErrSybAllowedShopExists) { + redirectSybShops(c, "", err.Error()) + return + } + fail(c, http.StatusInternalServerError, "新增同步店铺失败,已有店铺没有被改动。") + return + } + redirectSybShops(c, "店铺已加入允许列表并启用", "") +} + +func (h *Handler) SybAllowedShopStatus(c *gin.Context) { + enabled := c.PostForm("enabled") == "1" + err := service.SetSybAllowedShopEnabled(h.db, currentUser(c), c.PostForm("shop_id"), enabled, time.Now()) + if err != nil { + if service.IsValidationError(err) { + redirectSybShops(c, "", err.Error()) + return + } + fail(c, http.StatusInternalServerError, "更新同步店铺失败,原状态保持不变。") + return + } + message := "店铺已停用,后续同步会跳过该店铺" + if enabled { + message = "店铺已重新启用" + } + redirectSybShops(c, message, "") +} + +func redirectSybShops(c *gin.Context, message, errorMessage string) { + values := url.Values{} + if message != "" { + values.Set("msg", message) + } + if errorMessage != "" { + values.Set("error", errorMessage) + } + target := "/syb/shops" + if encoded := values.Encode(); encoded != "" { + target += "?" + encoded + } + c.Redirect(http.StatusSeeOther, target) +} diff --git a/admin/handler/web/web.go b/admin/handler/web/web.go index 99c510d..3055358 100644 --- a/admin/handler/web/web.go +++ b/admin/handler/web/web.go @@ -84,6 +84,10 @@ func Register(r *gin.Engine, db *sql.DB, onlineThreshold time.Duration) { pages.POST("/syb/match", h.SybMatch) pages.POST("/syb/create-task", h.SybCreateTask) pages.POST("/syb/delete", h.SybDelete) + sybShops := pages.Group("/syb/shops", AdminRequired()) + sybShops.GET("", h.SybAllowedShopList) + sybShops.POST("/create", h.SybAllowedShopCreate) + sybShops.POST("/status", h.SybAllowedShopStatus) // 4. 采集采购 pages.GET("/tasks", h.TaskList) diff --git a/admin/model/model.go b/admin/model/model.go index 743c329..1994c30 100644 --- a/admin/model/model.go +++ b/admin/model/model.go @@ -248,6 +248,9 @@ type SybSyncRun struct { DateTo string Status SybSyncRunStatus StockCount int + AcceptedCount int + ShopSkipped int + ShopFilterHash string DetailCount int Created int Updated int @@ -258,6 +261,17 @@ type SybSyncRun struct { FinishedAt string } +// SybAllowedShop 是顺运宝同步的全局店铺准入项。 +type SybAllowedShop struct { + ShopID string + ShopName string + NormalizedName string + Enabled bool + CreatedByUserID string + CreatedAt string + UpdatedAt string +} + // SpecMapping 是「这个蝦皮商品的顺运宝规格 = 拼多多的那个规格」。 // // # 为什么主键要带上 PddGoodsID diff --git a/admin/repository/mysql_db.go b/admin/repository/mysql_db.go index 92797e0..91c5c08 100644 --- a/admin/repository/mysql_db.go +++ b/admin/repository/mysql_db.go @@ -20,7 +20,7 @@ import ( "cmautobuy/admin/spec" ) -const mysqlSchemaVersion = 13 +const mysqlSchemaVersion = 14 // OpenMySQL 打开生产 MySQL 8 数据库。错误信息绝不包含完整 DSN 或密码。 func OpenMySQL(cfg config.DatabaseConfig) (*sql.DB, error) { @@ -580,10 +580,61 @@ func MigrateMySQL(db *sql.DB) error { if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 13, time.Now().UTC().Format(time.RFC3339Nano)); err != nil { return fmt.Errorf("记录 MySQL schema v13 失败: %w", err) } + current = 13 + } + if current < 14 { + if err := migrateMySQLV14(db); err != nil { + return fmt.Errorf("执行 MySQL schema v14 失败: %w", err) + } + if err := checkMySQLV14Shape(db); err != nil { + return fmt.Errorf("MySQL schema v14 自检失败,未记录版本: %w", err) + } + if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 14, time.Now().UTC().Format(time.RFC3339Nano)); err != nil { + return fmt.Errorf("记录 MySQL schema v14 失败: %w", err) + } } return CheckMySQLSchema(db) } +// migrateMySQLV14 增加顺运宝同步店铺准入表和审计统计。 +func migrateMySQLV14(db *sql.DB) error { + if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS syb_allowed_shops ( + shop_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY, + shop_name VARCHAR(500) NOT NULL, + normalized_name VARCHAR(500) COLLATE utf8mb4_bin NOT NULL, + enabled TINYINT NOT NULL DEFAULT 1, + created_by_user_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL, + created_at VARCHAR(35) NOT NULL, + updated_at VARCHAR(35) NOT NULL, + UNIQUE KEY uq_syb_allowed_shops_name (normalized_name), + KEY idx_syb_allowed_shops_enabled (enabled, normalized_name), + CONSTRAINT fk_syb_allowed_shops_user FOREIGN KEY (created_by_user_id) REFERENCES users(user_id), + CONSTRAINT chk_syb_allowed_shops_enabled CHECK (enabled IN (0,1)) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`); err != nil { + return fmt.Errorf("建立顺运宝允许店铺表失败: %w", err) + } + columns := []struct{ name, ddl string }{ + {"accepted_stock_count", `ALTER TABLE syb_sync_runs ADD COLUMN accepted_stock_count BIGINT NOT NULL DEFAULT 0 AFTER stock_count`}, + {"shop_skipped_count", `ALTER TABLE syb_sync_runs ADD COLUMN shop_skipped_count BIGINT NOT NULL DEFAULT 0 AFTER accepted_stock_count`}, + {"shop_filter_hash", `ALTER TABLE syb_sync_runs ADD COLUMN shop_filter_hash VARCHAR(64) COLLATE utf8mb4_bin NULL AFTER shop_skipped_count`}, + } + for _, column := range columns { + exists, err := mysqlColumnExists(db, "syb_sync_runs", column.name) + if err != nil { + return err + } + if !exists { + if _, err := db.Exec(column.ddl); err != nil { + return fmt.Errorf("增加 syb_sync_runs.%s 失败: %w", column.name, err) + } + } + } + if _, err := db.Exec(`UPDATE syb_sync_runs SET accepted_stock_count=stock_count, shop_skipped_count=0 WHERE shop_filter_hash IS NULL`); err != nil { + return fmt.Errorf("回填顺运宝同步店铺统计失败: %w", err) + } + return nil +} + // migrateMySQLV13 把任务真实主键迁为采集 cjN、采购 cgN 两套独立业务单号。 // // 外键 DDL 可重放;主键和关联数据在单一事务中同时切换。若进程在记录版本前退出, @@ -1490,6 +1541,7 @@ func CheckMySQLSchema(db *sql.DB) error { "catalog_import_runs", "task_syb_sources", "task_sequences", + "syb_allowed_shops", } if err := checkMySQLSchema(db, mysqlRequiredTables); err != nil { return err @@ -1524,7 +1576,51 @@ func CheckMySQLSchema(db *sql.DB) error { if err := checkMySQLV12Shape(db); err != nil { return err } - return checkMySQLV13Shape(db) + if err := checkMySQLV13Shape(db); err != nil { + return err + } + return checkMySQLV14Shape(db) +} + +func checkMySQLV14Shape(db *sql.DB) error { + if err := checkMySQLSchema(db, []string{"syb_allowed_shops"}); err != nil { + return err + } + for _, column := range []struct { + table, name string + length int64 + nullable bool + collation, defVal string + }{ + {"syb_allowed_shops", "shop_id", 191, false, "utf8mb4_bin", ""}, + {"syb_allowed_shops", "shop_name", 500, false, "utf8mb4_0900_ai_ci", ""}, + {"syb_allowed_shops", "normalized_name", 500, false, "utf8mb4_bin", ""}, + {"syb_allowed_shops", "created_by_user_id", 191, false, "utf8mb4_bin", ""}, + {"syb_sync_runs", "shop_filter_hash", 64, true, "utf8mb4_bin", ""}, + } { + if err := checkMySQLVarcharColumn(db, column.table, column.name, column.length, column.nullable, column.collation, column.defVal); err != nil { + return err + } + } + for _, name := range []string{"accepted_stock_count", "shop_skipped_count"} { + exists, err := mysqlColumnExists(db, "syb_sync_runs", name) + if err != nil || !exists { + return fmt.Errorf("顺运宝同步记录字段 %s 缺失", name) + } + } + for _, name := range []string{"uq_syb_allowed_shops_name", "idx_syb_allowed_shops_enabled"} { + exists, err := mysqlIndexExists(db, "syb_allowed_shops", name) + if err != nil || !exists { + return fmt.Errorf("顺运宝允许店铺索引 %s 缺失", name) + } + } + for _, name := range []string{"fk_syb_allowed_shops_user", "chk_syb_allowed_shops_enabled"} { + exists, err := mysqlConstraintExists(db, "syb_allowed_shops", name) + if err != nil || !exists { + return fmt.Errorf("顺运宝允许店铺约束 %s 缺失", name) + } + } + return nil } func checkMySQLV9Shape(db *sql.DB) error { diff --git a/admin/repository/mysql_db_integration_test.go b/admin/repository/mysql_db_integration_test.go index 8e5a270..49b1fa8 100644 --- a/admin/repository/mysql_db_integration_test.go +++ b/admin/repository/mysql_db_integration_test.go @@ -799,6 +799,42 @@ func TestNextTaskID_MySQL并发独立递增且回滚不耗号(t *testing.T) { } } +func TestMySQLMigrate_V13升级V14并可重放(t *testing.T) { + db := openMySQLMigrationTestDB(t) + defer db.Close() + cleanMySQLTestSchema(t, db) + defer cleanMySQLTestSchema(t, db) + if err := MigrateMySQL(db); err != nil { + t.Fatal(err) + } + + // 模拟生产 v13,并保留一条旧同步记录核对回填。 + mustExec(t, db, `DELETE FROM syb_allowed_shops`) + mustExec(t, db, `ALTER TABLE syb_sync_runs DROP COLUMN shop_filter_hash`) + mustExec(t, db, `ALTER TABLE syb_sync_runs DROP COLUMN shop_skipped_count`) + mustExec(t, db, `ALTER TABLE syb_sync_runs DROP COLUMN accepted_stock_count`) + mustExec(t, db, `DROP TABLE syb_allowed_shops`) + mustExec(t, db, `DELETE FROM schema_migrations WHERE version=14`) + mustExec(t, db, `INSERT INTO users(user_id,username,password_hash,role,status,password_changed_at,created_at,updated_at) + VALUES('V14-USER','v14-user','x','admin','active','2026-08-12T00:00:00Z','2026-08-12T00:00:00Z','2026-08-12T00:00:00Z')`) + mustExec(t, db, `INSERT INTO syb_sync_runs(run_id,user_id,date_from,date_to,status,stock_count,started_at,finished_at) + VALUES('V14-RUN','V14-USER','2026-08-11','2026-08-11','succeeded',7,'2026-08-12T00:00:00Z','2026-08-12T00:01:00Z')`) + + if err := MigrateMySQL(db); err != nil { + t.Fatalf("v13 升级 v14 失败: %v", err) + } + var accepted, skipped int + if err := db.QueryRow(`SELECT accepted_stock_count,shop_skipped_count FROM syb_sync_runs WHERE run_id='V14-RUN'`).Scan(&accepted, &skipped); err != nil || accepted != 7 || skipped != 0 { + t.Fatalf("旧同步记录回填错误: accepted=%d skipped=%d err=%v", accepted, skipped, err) + } + if err := MigrateMySQL(db); err != nil { + t.Fatalf("v14 重放失败: %v", err) + } + if err := checkMySQLV14Shape(db); err != nil { + t.Fatalf("v14 自检失败: %v", err) + } +} + func openMySQLMigrationTestDB(t *testing.T) *sql.DB { t.Helper() if os.Getenv("CMAUTOBUY_MYSQL_TEST") != "1" { diff --git a/admin/repository/sqlite_to_mysql.go b/admin/repository/sqlite_to_mysql.go index 308ab22..b1cea38 100644 --- a/admin/repository/sqlite_to_mysql.go +++ b/admin/repository/sqlite_to_mysql.go @@ -123,6 +123,12 @@ func MigrateSQLiteToMySQL(source, target *sql.DB, dryRun bool) (*SQLiteMigration if err := normalizeMigratedShopeeSKUs(tx); err != nil { return nil, err } + // 历史 SQLite 没有 v14 店铺准入统计;迁移后的旧记录按“未筛选”解释, + // 保持接受数等于原始货运单数。 + if _, err := tx.Exec(`UPDATE syb_sync_runs SET accepted_stock_count=stock_count, + shop_skipped_count=0,shop_filter_hash=NULL`); err != nil { + return nil, &SQLiteMigrationError{Stage: "无法补齐同步记录店铺统计", Table: "syb_sync_runs", Cause: err} + } if err := tx.Commit(); err != nil { return nil, &SQLiteMigrationError{Stage: "无法提交 MySQL 导入事务", Cause: err} } diff --git a/admin/repository/syb.go b/admin/repository/syb.go index dc5a1e2..d746e98 100644 --- a/admin/repository/syb.go +++ b/admin/repository/syb.go @@ -10,10 +10,86 @@ import ( "fmt" "strings" + "github.com/go-sql-driver/mysql" + "cmautobuy/admin/model" "cmautobuy/admin/spec" ) +var ErrSybAllowedShopExists = errors.New("顺运宝允许店铺已经存在") + +// ListSybAllowedShops 返回全部准入项;管理页需要同时看到已停用项。 +func ListSybAllowedShops(q Execer) ([]model.SybAllowedShop, error) { + rows, err := q.Query(`SELECT shop_id,shop_name,normalized_name,enabled,created_by_user_id,created_at,updated_at + FROM syb_allowed_shops ORDER BY enabled DESC,normalized_name`) + if err != nil { + return nil, fmt.Errorf("查询顺运宝允许店铺失败: %w", err) + } + defer rows.Close() + var result []model.SybAllowedShop + for rows.Next() { + var item model.SybAllowedShop + var enabled int + if err := rows.Scan(&item.ShopID, &item.ShopName, &item.NormalizedName, &enabled, + &item.CreatedByUserID, &item.CreatedAt, &item.UpdatedAt); err != nil { + return nil, fmt.Errorf("读取顺运宝允许店铺失败: %w", err) + } + item.Enabled = enabled == 1 + result = append(result, item) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("读取顺运宝允许店铺失败: %w", err) + } + return result, nil +} + +func ListEnabledSybShopNames(q Execer) ([]string, error) { + rows, err := q.Query(`SELECT normalized_name FROM syb_allowed_shops WHERE enabled=1 ORDER BY normalized_name`) + if err != nil { + return nil, fmt.Errorf("查询启用的顺运宝店铺失败: %w", err) + } + defer rows.Close() + var names []string + for rows.Next() { + var name string + if err := rows.Scan(&name); err != nil { + return nil, fmt.Errorf("读取启用的顺运宝店铺失败: %w", err) + } + names = append(names, name) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("读取启用的顺运宝店铺失败: %w", err) + } + return names, nil +} + +func InsertSybAllowedShop(q Execer, item model.SybAllowedShop) error { + _, err := q.Exec(`INSERT INTO syb_allowed_shops + (shop_id,shop_name,normalized_name,enabled,created_by_user_id,created_at,updated_at) + VALUES(?,?,?,?,?,?,?)`, item.ShopID, item.ShopName, item.NormalizedName, item.Enabled, + item.CreatedByUserID, item.CreatedAt, item.UpdatedAt) + if err != nil { + var mysqlErr *mysql.MySQLError + if errors.As(err, &mysqlErr) && mysqlErr.Number == 1062 { + return ErrSybAllowedShopExists + } + return fmt.Errorf("新增顺运宝允许店铺失败: %w", err) + } + return nil +} + +func SetSybAllowedShopEnabled(q Execer, shopID string, enabled bool, updatedAt string) (bool, error) { + result, err := q.Exec(`UPDATE syb_allowed_shops SET enabled=?,updated_at=? WHERE shop_id=?`, enabled, updatedAt, shopID) + if err != nil { + return false, fmt.Errorf("更新顺运宝允许店铺失败: %w", err) + } + affected, err := result.RowsAffected() + if err != nil { + return false, fmt.Errorf("确认顺运宝允许店铺更新结果失败: %w", err) + } + return affected == 1, nil +} + // ---------- 会话缓存 ---------- // SaveSybSession 写入或更新顺运宝登录会话缓存(按用户名 upsert)。 @@ -139,11 +215,13 @@ func FinishSybSyncRun(q Execer, run model.SybSyncRun) error { } result, err := q.Exec(` UPDATE syb_sync_runs - SET status = ?, stock_count = ?, detail_count = ?, created_count = ?, + SET status = ?, stock_count = ?, accepted_stock_count = ?, shop_skipped_count = ?, + shop_filter_hash = ?, detail_count = ?, created_count = ?, updated_count = ?, skipped_count = ?, error_message = ?, cursor_advanced = ?, finished_at = ? WHERE run_id = ? AND status = 'running'`, - run.Status, run.StockCount, run.DetailCount, run.Created, run.Updated, run.Skipped, + run.Status, run.StockCount, run.AcceptedCount, run.ShopSkipped, nullableText(run.ShopFilterHash), + run.DetailCount, run.Created, run.Updated, run.Skipped, nullableText(run.ErrorMessage), run.CursorAdvanced, run.FinishedAt, run.RunID) if err != nil { return fmt.Errorf("完成顺运宝同步记录失败: %w", err) @@ -179,7 +257,8 @@ func InterruptRunningSybSyncRuns(q Execer, finishedAt string) (int, error) { func ListSybSyncRuns(q Execer, limit, offset int) ([]model.SybSyncRun, error) { rows, err := q.Query(` SELECT r.run_id, r.user_id, u.username, r.date_from, r.date_to, r.status, - r.stock_count, r.detail_count, r.created_count, r.updated_count, + r.stock_count, r.accepted_stock_count, r.shop_skipped_count, r.shop_filter_hash, + r.detail_count, r.created_count, r.updated_count, r.skipped_count, r.error_message, r.cursor_advanced, r.started_at, r.finished_at FROM syb_sync_runs r @@ -194,14 +273,16 @@ func ListSybSyncRuns(q Execer, limit, offset int) ([]model.SybSyncRun, error) { var list []model.SybSyncRun for rows.Next() { var run model.SybSyncRun - var errorMessage, finishedAt sql.NullString + var errorMessage, finishedAt, shopFilterHash sql.NullString var cursorAdvanced int if err := rows.Scan(&run.RunID, &run.UserID, &run.Username, &run.DateFrom, &run.DateTo, - &run.Status, &run.StockCount, &run.DetailCount, &run.Created, &run.Updated, + &run.Status, &run.StockCount, &run.AcceptedCount, &run.ShopSkipped, &shopFilterHash, + &run.DetailCount, &run.Created, &run.Updated, &run.Skipped, &errorMessage, &cursorAdvanced, &run.StartedAt, &finishedAt); err != nil { return nil, fmt.Errorf("读取顺运宝同步记录失败: %w", err) } run.ErrorMessage = errorMessage.String + run.ShopFilterHash = shopFilterHash.String run.CursorAdvanced = cursorAdvanced == 1 run.FinishedAt = finishedAt.String list = append(list, run) diff --git a/admin/service/syb.go b/admin/service/syb.go index 7ec11b9..daed4b8 100644 --- a/admin/service/syb.go +++ b/admin/service/syb.go @@ -213,9 +213,12 @@ type SkipNote struct { // SyncReport 是一次同步的结果,供状态条显示。 type SyncReport struct { From, To string - Specified bool // true 表示操作员发起的指定日期补同步 - StockCount int // 拉到的货运单数 - DetailCount int // 落库的商品明细行数(不含跳过的) + Specified bool // true 表示操作员发起的指定日期补同步 + StockCount int // 拉到的货运单数 + AcceptedCount int // 店铺准入后接受的货运单数 + ShopSkipped int // 店铺不在允许列表或为空而跳过的货运单数 + ShopFilterHash string // 本次固定店铺快照的 SHA-256,不保存敏感凭据 + DetailCount int // 落库的商品明细行数(不含跳过的) Created int Updated int SkippedZero int // quantity <= 0 被跳过的条数 @@ -289,7 +292,9 @@ func FinishSybSyncRun(db *sql.DB, runID string, report SyncReport) error { } return repository.FinishSybSyncRun(db, model.SybSyncRun{ RunID: runID, Status: status, StockCount: report.StockCount, - DetailCount: report.DetailCount, Created: report.Created, Updated: report.Updated, + AcceptedCount: report.AcceptedCount, ShopSkipped: report.ShopSkipped, + ShopFilterHash: report.ShopFilterHash, + DetailCount: report.DetailCount, Created: report.Created, Updated: report.Updated, Skipped: report.SkippedZero, ErrorMessage: errorMessage, CursorAdvanced: report.CursorAdvanced, FinishedAt: finishedAt.UTC().Format(model.TimeLayout), }) @@ -324,8 +329,8 @@ func ListSybSyncHistory(db *sql.DB, page int) (*SybSyncHistoryResult, error) { RunID: run.RunID, Username: run.Username, DateRange: run.DateFrom + " ~ " + run.DateTo, StatusText: statusText, StatusClass: statusClass, - Summary: fmt.Sprintf("货运单 %d,明细 %d(新增 %d,更新 %d,跳过 %d)", - run.StockCount, run.DetailCount, run.Created, run.Updated, run.Skipped), + Summary: fmt.Sprintf("原始货运单 %d,接受 %d,店铺跳过 %d;明细 %d(新增 %d,更新 %d,数量跳过 %d)", + run.StockCount, run.AcceptedCount, run.ShopSkipped, run.DetailCount, run.Created, run.Updated, run.Skipped), ErrorMessage: run.ErrorMessage, StartedAt: formatLocalTime(run.StartedAt), FinishedAt: formatLocalTime(run.FinishedAt), } @@ -362,8 +367,8 @@ func (r SyncReport) Summary() string { if r.Err != nil { return "同步失败:" + r.Err.Error() } - msg := fmt.Sprintf("同步完成:日期范围 %s ~ %s,货运单 %d 张,商品明细 %d 条(新增 %d,更新 %d,跳过 %d)", - r.From, r.To, r.StockCount, r.DetailCount, r.Created, r.Updated, r.SkippedZero) + msg := fmt.Sprintf("同步完成:日期范围 %s ~ %s,原始货运单 %d 张,接受 %d 张,店铺跳过 %d 张;商品明细 %d 条(新增 %d,更新 %d,数量跳过 %d)", + r.From, r.To, r.StockCount, r.AcceptedCount, r.ShopSkipped, r.DetailCount, r.Created, r.Updated, r.SkippedZero) if len(r.Notes) > 0 { var reasons []string for _, n := range r.Notes { @@ -548,6 +553,22 @@ func RunSybSync(ctx context.Context, db *sql.DB, client *syb.Client, cfg config. // 局部历史补拉只 upsert 数据、不动游标,否则会让未覆盖的订单永久漏掉。 func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client, cfg config.SybConfig, now time.Time, options SybSyncOptions) SyncReport { report := SyncReport{StartedAt: now, Specified: options.IsSpecified()} + allowedNames, err := repository.ListEnabledSybShopNames(db) + if err != nil { + report.Err = fmt.Errorf("读取顺运宝允许店铺失败: %w", err) + report.FinishedAt = time.Now().UTC() + return report + } + if len(allowedNames) == 0 { + report.Err = fmt.Errorf("没有启用的顺运宝同步店铺,请先由管理员在“同步店铺”中配置并启用至少一个店铺") + report.FinishedAt = time.Now().UTC() + return report + } + allowedShops := make(map[string]struct{}, len(allowedNames)) + for _, name := range allowedNames { + allowedShops[name] = struct{}{} + } + report.ShopFilterHash = fmt.Sprintf("%x", sha256.Sum256([]byte(strings.Join(allowedNames, "\x00")))) pageSize := cfg.PageSize if pageSize <= 0 { @@ -729,8 +750,17 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client, } stockByID := listResult.stockByID - orderedIDs := listResult.orderedIDs - report.StockCount += len(orderedIDs) + rawIDs := listResult.orderedIDs + report.StockCount += len(rawIDs) + orderedIDs := make([]int64, 0, len(rawIDs)) + for _, id := range rawIDs { + if !sybShopAllowed(allowedShops, stockByID[id].Raw, nil) { + report.ShopSkipped++ + continue + } + orderedIDs = append(orderedIDs, id) + } + report.AcceptedCount += len(orderedIDs) for i := 0; i < len(orderedIDs); i += detailBatch { end := i + detailBatch @@ -752,6 +782,11 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client, } for _, d := range details { stockRow := stockByID[d.ID] + if !sybShopAllowed(allowedShops, stockRow.Raw, d.Raw) { + report.AcceptedCount-- + report.ShopSkipped++ + continue + } if err := writeStockDetail(db, cfg.BaseURL, stockRow, d, &report); err != nil { report.Err = fmt.Errorf("写入货运单 %s(id=%d)失败(本次同步整体作废,"+ "已写入的数据保留): %w", d.Code, d.ID, err) @@ -762,9 +797,9 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client, } if unstableTodayErr != nil { report.Err = fmt.Errorf( - "%s 当天货运单在连续 %d 次分页期间仍有变化;已保存最后一次取得的 %d 张货运单完整明细,"+ + "%s 当天货运单在连续 %d 次分页期间仍有变化;最后一次取得原始货运单 %d 张,已保存其中允许店铺 %d 张的完整明细,"+ "本次未形成稳定快照且不推进游标,下次同步将继续覆盖当天:%w", - plan.date, sybTodayListMaxAttempts, len(orderedIDs), unstableTodayErr) + plan.date, sybTodayListMaxAttempts, len(rawIDs), len(orderedIDs), unstableTodayErr) report.FinishedAt = time.Now().UTC() return report } @@ -783,6 +818,14 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client, return report } +// sybShopAllowed 用明细字段覆盖列表字段后再核对,防止列表通过但明细在同步 +// 期间已变成其他店铺。detail 为空时只检查列表快照。 +func sybShopAllowed(allowed map[string]struct{}, listRaw, detailRaw map[string]any) bool { + name := trimmedStringField(mergeRaw(listRaw, detailRaw), "shopName") + _, ok := allowed[name] + return ok +} + // validateDetailBatch 确认批量明细响应与请求 ID 一一对应。任何缺失、重复、 // 意外 ID 或空商品明细都会让同步失败,避免在数据不完整时推进游标。 func validateDetailBatch(requested []int64, details []syb.StockDetail) error { diff --git a/admin/service/syb_allowed_shop.go b/admin/service/syb_allowed_shop.go new file mode 100644 index 0000000..eeeb803 --- /dev/null +++ b/admin/service/syb_allowed_shop.go @@ -0,0 +1,80 @@ +package service + +import ( + "database/sql" + "errors" + "fmt" + "strings" + "time" + + "cmautobuy/admin/model" + "cmautobuy/admin/repository" +) + +const maxSybShopNameLength = 500 + +func ListSybAllowedShops(db *sql.DB, actor *model.User) ([]model.SybAllowedShop, error) { + if actor == nil || !actor.IsAdmin() { + return nil, ErrAdminRequired + } + return repository.ListSybAllowedShops(db) +} + +func CountEnabledSybAllowedShops(db *sql.DB) (int, error) { + names, err := repository.ListEnabledSybShopNames(db) + return len(names), err +} + +func EnsureEnabledSybAllowedShops(db *sql.DB) error { + count, err := CountEnabledSybAllowedShops(db) + if err != nil { + return err + } + if count == 0 { + return &validationError{field: "shop_name", message: "没有启用的顺运宝同步店铺,请先由管理员配置并启用至少一个店铺"} + } + return nil +} + +func CreateSybAllowedShop(db *sql.DB, actor *model.User, rawName string, now time.Time) error { + if actor == nil || !actor.IsAdmin() { + return ErrAdminRequired + } + name := strings.TrimSpace(rawName) + if name == "" { + return &validationError{field: "shop_name", message: "店铺名称不能为空"} + } + if len([]rune(name)) > maxSybShopNameLength { + return &validationError{field: "shop_name", message: "店铺名称最多 500 个字符"} + } + id, err := randomID("SHOP-", 16) + if err != nil { + return fmt.Errorf("生成店铺编号失败: %w", err) + } + at := now.UTC().Format(model.TimeLayout) + err = repository.InsertSybAllowedShop(db, model.SybAllowedShop{ + ShopID: id, ShopName: name, NormalizedName: name, Enabled: true, + CreatedByUserID: actor.UserID, CreatedAt: at, UpdatedAt: at, + }) + if errors.Is(err, repository.ErrSybAllowedShopExists) { + return &validationError{field: "shop_name", message: "该店铺已经在允许列表中,可直接重新启用"} + } + return err +} + +func SetSybAllowedShopEnabled(db *sql.DB, actor *model.User, shopID string, enabled bool, now time.Time) error { + if actor == nil || !actor.IsAdmin() { + return ErrAdminRequired + } + if strings.TrimSpace(shopID) == "" { + return &validationError{field: "shop_id", message: "店铺编号不能为空"} + } + found, err := repository.SetSybAllowedShopEnabled(db, shopID, enabled, now.UTC().Format(model.TimeLayout)) + if err != nil { + return err + } + if !found { + return &validationError{field: "shop_id", message: "店铺不存在,请刷新页面后重试"} + } + return nil +} diff --git a/admin/service/syb_allowed_shop_test.go b/admin/service/syb_allowed_shop_test.go new file mode 100644 index 0000000..78928d5 --- /dev/null +++ b/admin/service/syb_allowed_shop_test.go @@ -0,0 +1,51 @@ +package service + +import ( + "errors" + "testing" + "time" + + "cmautobuy/admin/model" + "cmautobuy/admin/repository" +) + +func TestSybAllowedShop_管理员维护与精确去重(t *testing.T) { + db := newTestDB(t) + now := time.Date(2026, 8, 12, 10, 0, 0, 0, time.UTC) + admin := &model.User{UserID: "SHOP-ADMIN", Username: "shop-admin", PasswordHash: "x", + Role: model.RoleAdmin, Status: model.UserActive, PasswordChangedAt: model.NowISO(), + CreatedAt: model.NowISO(), UpdatedAt: model.NowISO()} + if err := repository.CreateUser(db, *admin); err != nil { + t.Fatal(err) + } + + if err := CreateSybAllowedShop(db, admin, " qwg8fkb044 ", now); err != nil { + t.Fatal(err) + } + rows, err := ListSybAllowedShops(db, admin) + if err != nil || len(rows) != 1 || rows[0].ShopName != "qwg8fkb044" || !rows[0].Enabled { + t.Fatalf("新增结果错误: rows=%+v err=%v", rows, err) + } + if err := CreateSybAllowedShop(db, admin, "qwg8fkb044", now); !IsValidationError(err) { + t.Fatalf("去除首尾空白后的重名应是表单错误,实际 %v", err) + } + if err := SetSybAllowedShopEnabled(db, admin, rows[0].ShopID, false, now.Add(time.Minute)); err != nil { + t.Fatal(err) + } + if count, err := CountEnabledSybAllowedShops(db); err != nil || count != 0 { + t.Fatalf("停用后启用数错误: count=%d err=%v", count, err) + } + if err := EnsureEnabledSybAllowedShops(db); !IsValidationError(err) { + t.Fatalf("空白名单应阻止同步: %v", err) + } +} + +func TestSybAllowedShop_采购员不能管理(t *testing.T) { + purchaser := &model.User{UserID: "BUYER", Role: model.RolePurchaser, Status: model.UserActive} + if err := CreateSybAllowedShop(nil, purchaser, "店铺", time.Now()); !errors.Is(err, ErrAdminRequired) { + t.Fatalf("采购员新增应被拒绝: %v", err) + } + if _, err := ListSybAllowedShops(nil, purchaser); !errors.Is(err, ErrAdminRequired) { + t.Fatalf("采购员读取管理列表应被拒绝: %v", err) + } +} diff --git a/admin/service/syb_test.go b/admin/service/syb_test.go index 1c2f84f..d51ff25 100644 --- a/admin/service/syb_test.go +++ b/admin/service/syb_test.go @@ -5,8 +5,11 @@ import ( "database/sql" "encoding/json" "fmt" + "io" "net/http" "net/http/httptest" + "net/http/httputil" + "net/url" "strings" "testing" "time" @@ -198,6 +201,9 @@ func fakeSybServer(t *testing.T, stocks []fakeStock, failListPageIndex int) *htt if stocks[i].Created == "" { stocks[i].Created = "2026-07-28" } + if stocks[i].ShopName == "" { + stocks[i].ShopName = "测试店铺" + } s := stocks[i] byID[s.ID] = s } @@ -315,6 +321,16 @@ type integrityServerData struct { func fakeIntegritySybServer(t *testing.T, data integrityServerData) *httptest.Server { t.Helper() + for _, row := range data.list { + if _, ok := row["shopName"]; !ok { + row["shopName"] = "测试店铺" + } + } + for _, row := range data.details { + if _, ok := row["shopName"]; !ok { + row["shopName"] = "测试店铺" + } + } totalCalls := 0 return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch r.URL.Path { @@ -349,7 +365,100 @@ func writeEnvelope(t *testing.T, w http.ResponseWriter, status bool, msg string, func newSyncTestDB(t *testing.T) *sql.DB { t.Helper() - return newTestDB(t) + db := newTestDB(t) + now := model.NowISO() + if err := repository.CreateUser(db, model.User{ + UserID: "SYB-TEST-ADMIN", Username: "syb-test-admin", PasswordHash: "test", + Role: model.RoleAdmin, Status: model.UserActive, PasswordChangedAt: now, CreatedAt: now, UpdatedAt: now, + }); err != nil { + t.Fatalf("准备同步测试管理员失败: %v", err) + } + if err := repository.InsertSybAllowedShop(db, model.SybAllowedShop{ + ShopID: "SYB-TEST-SHOP", ShopName: "测试店铺", NormalizedName: "测试店铺", Enabled: true, + CreatedByUserID: "SYB-TEST-ADMIN", CreatedAt: now, UpdatedAt: now, + }); err != nil { + t.Fatalf("准备同步测试允许店铺失败: %v", err) + } + return db +} + +func TestRunSybSync_没有启用店铺时不请求顺运宝(t *testing.T) { + db := newSyncTestDB(t) + if _, err := db.Exec(`UPDATE syb_allowed_shops SET enabled=0`); err != nil { + t.Fatal(err) + } + requests := 0 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + requests++ + w.WriteHeader(http.StatusInternalServerError) + })) + defer srv.Close() + client, _ := syb.New(srv.URL) + report := RunSybSyncWithOptions(context.Background(), db, client, + config.SybConfig{BaseURL: srv.URL, PageSize: 20, MaxMatches: 100}, time.Now(), + SybSyncOptions{From: "2026-08-11", To: "2026-08-11"}) + if report.Err == nil || !strings.Contains(report.Err.Error(), "没有启用") { + t.Fatalf("实际错误: %v", report.Err) + } + if requests != 0 { + t.Fatalf("空白名单不应请求顺运宝,实际 %d 次", requests) + } + if report.CursorAdvanced { + t.Fatal("空白名单不得推进游标") + } +} + +func TestSybShopAllowed_明细店铺覆盖列表后重新拦截(t *testing.T) { + allowed := map[string]struct{}{"测试店铺": {}} + if !sybShopAllowed(allowed, map[string]any{"shopName": " 测试店铺 "}, nil) { + t.Fatal("应忽略允许店铺名称首尾空白") + } + if sybShopAllowed(allowed, map[string]any{"shopName": "测试店铺"}, map[string]any{"shopName": "其他店铺"}) { + t.Fatal("明细店铺变化后必须重新拦截") + } + if sybShopAllowed(allowed, map[string]any{"shopName": "测试店铺"}, map[string]any{"shopName": " "}) { + t.Fatal("明细店铺变为空值时必须拦截") + } +} + +func TestRunSybSync_只请求并写入允许店铺(t *testing.T) { + db := newSyncTestDB(t) + detailIDs := []int64{} + stocks := []fakeStock{ + {ID: 1, Code: "A", ShopName: " 测试店铺 ", Details: []fakeDetail{{ID: 11, ProductID: 111, ProductQty: 1}}}, + {ID: 2, Code: "B", ShopName: "其他店铺", Details: []fakeDetail{{ID: 22, ProductID: 222, ProductQty: 1}}}, + } + base := fakeSybServer(t, stocks, 0) + defer base.Close() + target, _ := url.Parse(base.URL) + reverseProxy := httputil.NewSingleHostReverseProxy(target) + proxy := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/am/stock/detail/listByStock" { + var body struct { + IDs []int64 `json:"ids"` + } + data, _ := io.ReadAll(r.Body) + _ = json.Unmarshal(data, &body) + detailIDs = append(detailIDs, body.IDs...) + r.Body = io.NopCloser(strings.NewReader(string(data))) + } + reverseProxy.ServeHTTP(w, r) + })) + defer proxy.Close() + client, _ := syb.New(proxy.URL) + report := RunSybSyncWithOptions(context.Background(), db, client, + config.SybConfig{BaseURL: proxy.URL, PageSize: 20, MaxMatches: 100}, + time.Date(2026, 7, 28, 12, 0, 0, 0, time.UTC), + SybSyncOptions{From: "2026-07-28", To: "2026-07-28"}) + if report.Err != nil { + t.Fatal(report.Err) + } + if report.StockCount != 2 || report.AcceptedCount != 1 || report.ShopSkipped != 1 || report.Created != 1 { + t.Fatalf("店铺统计不正确: %+v", report) + } + if len(detailIDs) != 1 || detailIDs[0] != 1 { + t.Fatalf("明细请求应只有允许店铺,实际 %v", detailIDs) + } } func TestRunSybSync_已有的ShopeeSKUID同步后仍在(t *testing.T) { @@ -789,7 +898,7 @@ func fakeGrowingTodaySybServer(t *testing.T, today string, growAttempts int) (*h rows := make([]map[string]any, 0, count) for i := 1; i <= count; i++ { id := idOffset + int64(i) - rows = append(rows, map[string]any{"id": id, "code": fmt.Sprintf("ORDER-%d", id)}) + rows = append(rows, map[string]any{"id": id, "code": fmt.Sprintf("ORDER-%d", id), "shopName": "测试店铺"}) } return rows } @@ -833,7 +942,7 @@ func fakeGrowingTodaySybServer(t *testing.T, today string, growAttempts int) (*h list := make([]map[string]any, 0, len(body.IDs)) for _, id := range body.IDs { list = append(list, map[string]any{ - "id": id, "code": fmt.Sprintf("ORDER-%d", id), + "id": id, "code": fmt.Sprintf("ORDER-%d", id), "shopName": "测试店铺", "details": []map[string]any{{ "id": id + 1000, "productId": id + 10000, "productQty": 1, }}, diff --git a/admin/templates/syb/list.html b/admin/templates/syb/list.html index bb0d48d..636c255 100644 --- a/admin/templates/syb/list.html +++ b/admin/templates/syb/list.html @@ -21,6 +21,7 @@ + {{if .CurrentUser.IsAdmin}}同步店铺({{.EnabledSyncShopCount}}){{end}} @@ -66,7 +67,7 @@ {{if .RangeError}}{{end}} {{if .RangeWarning}}

{{.RangeWarning}}

{{end}}

- 按顺运宝货运单创建日期(UTC+8)同步,开始日和结束日都包含,每次最多 31 天。 + 按顺运宝货运单创建日期(UTC+8)同步,开始日和结束日都包含,每次最多 31 天;当前启用 {{.EnabledSyncShopCount}} 个同步店铺。