feat: 按允许店铺筛选顺运宝同步 (#196)
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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" {
|
||||
|
||||
@@ -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}
|
||||
}
|
||||
|
||||
+86
-5
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user