feat: 统一顺运宝日期同步并记录历史 (#59)

This commit is contained in:
chengma
2026-08-09 18:22:58 +08:00
parent ea7932db63
commit 05a54b044d
18 changed files with 934 additions and 144 deletions
+39 -4
View File
@@ -264,7 +264,7 @@ var migrations = [][]string{
// 背景见 #20:v1 曾经被原地改写而不是新增版本,导致已经建过库的机器
// (user_version 已经越过 v1)永远不会重跑改写后的语句,程序拿着一个
// 和代码对不上的库静默启动。
const schemaVersion = 7
const schemaVersion = 8
// migrationV4 给 PDD 商品增加店铺名。
//
@@ -353,6 +353,32 @@ var migrationV7 = []string{
ON client_user_assignments(client_id, started_at DESC);`,
}
// migrationV8 持久化顺运宝同步记录,见工单 #59。
// 只新增表和索引,不修改 v1-v7 的任何已发布语句。
var migrationV8 = []string{
`CREATE TABLE syb_sync_runs (
run_id TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
date_from TEXT NOT NULL,
date_to TEXT NOT NULL,
status TEXT NOT NULL CHECK (status IN ('running', 'succeeded', 'failed', 'interrupted')),
stock_count INTEGER NOT NULL DEFAULT 0,
detail_count INTEGER NOT NULL DEFAULT 0,
created_count INTEGER NOT NULL DEFAULT 0,
updated_count INTEGER NOT NULL DEFAULT 0,
skipped_count INTEGER NOT NULL DEFAULT 0,
error_message TEXT,
cursor_advanced INTEGER NOT NULL DEFAULT 0 CHECK (cursor_advanced IN (0, 1)),
started_at TEXT NOT NULL,
finished_at TEXT,
FOREIGN KEY (user_id) REFERENCES users(user_id),
CHECK ((status = 'running' AND finished_at IS NULL)
OR (status <> 'running' AND finished_at IS NOT NULL))
);`,
`CREATE INDEX idx_syb_sync_runs_started
ON syb_sync_runs(started_at DESC, run_id DESC);`,
}
// Migrate 把数据库升到最新版本。
// 已经是最新的就什么都不做,可以重复调用。
func Migrate(db *sql.DB) error {
@@ -431,6 +457,14 @@ func Migrate(db *sql.DB) error {
if err := runSQLMigration(db, 7, migrationV7); err != nil {
return err
}
reached = 7
}
// v8 只新增顺运宝同步记录表和倒序索引。
if reached < 8 {
if err := runSQLMigration(db, 8, migrationV8); err != nil {
return err
}
}
return nil
@@ -892,7 +926,7 @@ var requiredTables = []string{
"syb_orders", "sku_mappings", "tasks", "clients",
"idempotency_keys", "task_claims",
"syb_session", "syb_sync_state",
"users", "web_sessions", "client_user_assignments",
"users", "web_sessions", "client_user_assignments", "syb_sync_runs",
}
// requiredColumns 只列出不能靠“表存在”发现的关键追加列。
@@ -900,8 +934,9 @@ var requiredTables = []string{
// 查询对应页面会直接失败,因此启动时就应给出明确错误,
// 而不是等操作员点到页面才暴露。
var requiredColumns = map[string][]string{
"pdd_products": {"shop_name"},
"syb_orders": {"product_spec"},
"pdd_products": {"shop_name"},
"syb_orders": {"product_spec"},
"syb_sync_runs": {"user_id", "date_from", "date_to", "status", "cursor_advanced", "started_at", "finished_at"},
}
// CheckSchema 在 Migrate 成功后调用,确认代码依赖的表都在。
+66 -6
View File
@@ -9,6 +9,8 @@ import (
"sort"
"strings"
"testing"
"cmautobuy/admin/model"
)
// 本文件是工单 #20 的核心交付物:证明不管从哪种库起步,迁到最新版本后
@@ -623,11 +625,10 @@ func TestMigrate_v6用户表与v7客户端归属表均存在(t *testing.T) {
}
}
// #52/#54:显式覆盖每一个已发布版本起点,证明 Client 现实库
// 不会因为 v6/v7 的新增表而卡在中间版本。v2 的两种历史结构另由上面的收敛
// 测试持续覆盖;这里验证顺序发布的 v1-v5 路径。
func TestMigrate_v1到v5均可升级到最新版本(t *testing.T) {
for version := 1; version <= 5; version++ {
// 显式覆盖每一个已发布版本起点,证明现实库不会因为后续新增表而卡住。
// v2 的两种历史结构另由上面的收敛测试持续覆盖。
func TestMigrate_v1到v7均可升级到最新版本(t *testing.T) {
for version := 1; version <= 7; version++ {
t.Run(fmt.Sprintf("v%d", version), func(t *testing.T) {
db := newPublishedVersionDB(t, version)
if err := Migrate(db); err != nil {
@@ -647,6 +648,41 @@ func TestMigrate_v1到v5均可升级到最新版本(t *testing.T) {
}
}
func TestMigrate_v8新增顺运宝同步记录表(t *testing.T) {
db := newPublishedVersionDB(t, 7)
if err := Migrate(db); err != nil {
t.Fatalf("v7 升级到 v8 失败: %v", err)
}
if !existingTableSet(t, db)["syb_sync_runs"] {
t.Fatal("v8 应该新增 syb_sync_runs 表")
}
var version int
if err := db.QueryRow(`PRAGMA user_version`).Scan(&version); err != nil || version != 8 {
t.Fatalf("迁移版本错误: version=%d err=%v", version, err)
}
user := model.User{
UserID: "USR-MIGRATE", Username: "buyer", PasswordHash: "test-hash",
Role: model.RolePurchaser, Status: model.UserActive,
PasswordChangedAt: model.NowISO(), CreatedAt: model.NowISO(), UpdatedAt: model.NowISO(),
}
if err := CreateInitialAdmin(db, user); err != nil {
t.Fatalf("准备同步记录外键用户失败: %v", err)
}
// 状态和 running/finished_at 对应关系必须由数据库兜底,不能只靠 Go 校验。
if _, err := db.Exec(`INSERT INTO syb_sync_runs
(run_id, user_id, date_from, date_to, status, started_at, finished_at)
VALUES ('bad', ?, '2026-08-09', '2026-08-09', 'unknown', ?, ?)`,
user.UserID, model.NowISO(), model.NowISO()); err == nil {
t.Fatal("非法同步状态应该被数据库约束拒绝")
}
if _, err := db.Exec(`INSERT INTO syb_sync_runs
(run_id, user_id, date_from, date_to, status, started_at, finished_at)
VALUES ('bad-running', ?, '2026-08-09', '2026-08-09', 'running', ?, ?)`,
user.UserID, model.NowISO(), model.NowISO()); err == nil {
t.Fatal("running 状态不应允许填写完成时间")
}
}
func newPublishedVersionDB(t *testing.T, version int) *sql.DB {
t.Helper()
db, err := Open(t.TempDir())
@@ -683,6 +719,16 @@ func newPublishedVersionDB(t *testing.T, version int) *sql.DB {
t.Fatalf("构造 v%d 时执行 v5 失败: %v", version, err)
}
}
if version >= 6 {
if err := runSQLMigration(db, 6, migrationV6); err != nil {
t.Fatalf("构造 v%d 时执行 v6 失败: %v", version, err)
}
}
if version >= 7 {
if err := runSQLMigration(db, 7, migrationV7); err != nil {
t.Fatalf("构造 v%d 时执行 v7 失败: %v", version, err)
}
}
return db
}
@@ -1120,7 +1166,7 @@ func TestCheckSchema_缺少v4关键列时拒绝(t *testing.T) {
if err := migrateV3(db); err != nil {
t.Fatalf("准备 v3 数据库失败: %v", err)
}
// 故意跳过 v4(不加 shop_name),但把 v5-v7 补上——否则 CheckSchema 会先
// 故意跳过 v4(不加 shop_name),但把 v5-v8 补上——否则 CheckSchema 会先
// 因为缺后续表报错,测不到本测试真正要覆盖的"缺 shop_name"这条路径。
if err := runSQLMigration(db, 5, migrationV5); err != nil {
t.Fatalf("准备 v5 数据库失败: %v", err)
@@ -1131,6 +1177,9 @@ func TestCheckSchema_缺少v4关键列时拒绝(t *testing.T) {
if err := runSQLMigration(db, 7, migrationV7); err != nil {
t.Fatalf("准备 v7 数据库失败: %v", err)
}
if err := runSQLMigration(db, 8, migrationV8); err != nil {
t.Fatalf("准备 v8 数据库失败: %v", err)
}
err := CheckSchema(db)
if err == nil || !strings.Contains(err.Error(), "shop_name") {
@@ -1152,3 +1201,14 @@ func TestCheckSchema_缺表时拒绝(t *testing.T) {
t.Errorf("错误信息应该指出缺的是哪张表,实际: %v", err)
}
}
func TestCheckSchema_缺少同步记录表时拒绝(t *testing.T) {
db := newFreshDB(t)
if _, err := db.Exec(`DROP TABLE syb_sync_runs`); err != nil {
t.Fatalf("删表失败: %v", err)
}
err := CheckSchema(db)
if err == nil || !strings.Contains(err.Error(), "syb_sync_runs") {
t.Fatalf("缺少同步记录表时应拒绝启动并指出表名,实际 %v", err)
}
}
+107
View File
@@ -112,6 +112,113 @@ func SetSybLastSyncedAt(q Execer, at string) error {
return nil
}
// ---------- 同步记录 ----------
// CreateSybSyncRun 在真正启动后台同步前写入一条 running 记录。
func CreateSybSyncRun(q Execer, run model.SybSyncRun) error {
if run.RunID == "" || run.UserID == "" || run.DateFrom == "" || run.DateTo == "" || run.StartedAt == "" {
return fmt.Errorf("同步记录缺少编号、操作人、日期范围或开始时间")
}
_, err := q.Exec(`
INSERT INTO syb_sync_runs
(run_id, user_id, date_from, date_to, status, started_at)
VALUES (?, ?, ?, ?, 'running', ?)`,
run.RunID, run.UserID, run.DateFrom, run.DateTo, run.StartedAt)
if err != nil {
return fmt.Errorf("创建顺运宝同步记录失败: %w", err)
}
return nil
}
// FinishSybSyncRun 把 running 记录更新为最终状态。
func FinishSybSyncRun(q Execer, run model.SybSyncRun) error {
if run.Status != model.SybSyncSucceeded && run.Status != model.SybSyncFailed {
return fmt.Errorf("同步完成状态不合法: %s", run.Status)
}
result, err := q.Exec(`
UPDATE syb_sync_runs
SET status = ?, stock_count = ?, 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,
nullableText(run.ErrorMessage), run.CursorAdvanced, run.FinishedAt, run.RunID)
if err != nil {
return fmt.Errorf("完成顺运宝同步记录失败: %w", err)
}
affected, err := result.RowsAffected()
if err != nil {
return fmt.Errorf("确认顺运宝同步记录完成结果失败: %w", err)
}
if affected != 1 {
return fmt.Errorf("同步记录 %s 不存在或已经结束", run.RunID)
}
return nil
}
// InterruptRunningSybSyncRuns 在 Admin 启动时收敛上次进程遗留的 running 记录。
func InterruptRunningSybSyncRuns(q Execer, finishedAt string) (int, error) {
result, err := q.Exec(`
UPDATE syb_sync_runs
SET status = 'interrupted', finished_at = ?,
error_message = 'Admin 在同步完成前退出,请重新同步该日期范围'
WHERE status = 'running'`, finishedAt)
if err != nil {
return 0, fmt.Errorf("标记中断的顺运宝同步记录失败: %w", err)
}
affected, err := result.RowsAffected()
if err != nil {
return 0, fmt.Errorf("统计中断的顺运宝同步记录失败: %w", err)
}
return int(affected), nil
}
// ListSybSyncRuns 按开始时间倒序分页查询同步记录。
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.skipped_count, r.error_message, r.cursor_advanced,
r.started_at, r.finished_at
FROM syb_sync_runs r
JOIN users u ON u.user_id = r.user_id
ORDER BY r.started_at DESC, r.run_id DESC
LIMIT ? OFFSET ?`, limit, offset)
if err != nil {
return nil, fmt.Errorf("查询顺运宝同步记录失败: %w", err)
}
defer rows.Close()
var list []model.SybSyncRun
for rows.Next() {
var run model.SybSyncRun
var errorMessage, finishedAt 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.Skipped, &errorMessage, &cursorAdvanced, &run.StartedAt, &finishedAt); err != nil {
return nil, fmt.Errorf("读取顺运宝同步记录失败: %w", err)
}
run.ErrorMessage = errorMessage.String
run.CursorAdvanced = cursorAdvanced == 1
run.FinishedAt = finishedAt.String
list = append(list, run)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("读取顺运宝同步记录失败: %w", err)
}
return list, nil
}
// CountSybSyncRuns 返回同步记录总数,供弹窗分页。
func CountSybSyncRuns(q Execer) (int, error) {
var count int
if err := q.QueryRow(`SELECT COUNT(*) FROM syb_sync_runs`).Scan(&count); err != nil {
return 0, fmt.Errorf("统计顺运宝同步记录失败: %w", err)
}
return count, nil
}
// ---------- 货运单明细 ----------
// UpsertSybOrder 写入或更新一条顺运宝货运单明细行。
+76
View File
@@ -191,6 +191,82 @@ func TestSybSyncState_首次为空之后可更新(t *testing.T) {
}
}
func createSybSyncTestUser(t *testing.T, db *sql.DB) model.User {
t.Helper()
user := model.User{
UserID: "USR-SYNC", Username: "buyer", PasswordHash: "test-hash",
Role: model.RolePurchaser, Status: model.UserActive,
PasswordChangedAt: "2026-08-09T00:00:00Z", CreatedAt: "2026-08-09T00:00:00Z",
UpdatedAt: "2026-08-09T00:00:00Z",
}
if err := CreateInitialAdmin(db, user); err != nil {
t.Fatalf("创建同步记录测试用户失败: %v", err)
}
return user
}
func TestSybSyncRun_创建完成并分页读取(t *testing.T) {
db := newSybTestDB(t)
user := createSybSyncTestUser(t, db)
run := model.SybSyncRun{
RunID: "SYB-RUN-1", UserID: user.UserID, DateFrom: "2026-08-07",
DateTo: "2026-08-09", StartedAt: "2026-08-09T01:00:00Z",
}
if err := CreateSybSyncRun(db, run); err != nil {
t.Fatalf("创建同步记录失败: %v", err)
}
if err := FinishSybSyncRun(db, model.SybSyncRun{
RunID: "SYB-RUN-1", Status: model.SybSyncSucceeded,
StockCount: 3, DetailCount: 4, Created: 2, Updated: 2, Skipped: 1,
CursorAdvanced: true, FinishedAt: "2026-08-09T01:02:00Z",
}); err != nil {
t.Fatalf("完成同步记录失败: %v", err)
}
rows, err := ListSybSyncRuns(db, 10, 0)
if err != nil || len(rows) != 1 {
t.Fatalf("读取同步记录失败: rows=%+v err=%v", rows, err)
}
got := rows[0]
if got.Username != "buyer" || got.Status != model.SybSyncSucceeded ||
got.DetailCount != 4 || !got.CursorAdvanced || got.FinishedAt == "" {
t.Fatalf("同步记录字段不正确: %+v", got)
}
if count, err := CountSybSyncRuns(db); err != nil || count != 1 {
t.Fatalf("同步记录总数错误: count=%d err=%v", count, err)
}
}
func TestInterruptRunningSybSyncRuns_只中断未完成记录(t *testing.T) {
db := newSybTestDB(t)
user := createSybSyncTestUser(t, db)
for _, id := range []string{"RUNNING", "DONE"} {
if err := CreateSybSyncRun(db, model.SybSyncRun{
RunID: id, UserID: user.UserID, DateFrom: "2026-08-09", DateTo: "2026-08-09",
StartedAt: "2026-08-09T01:00:00Z",
}); err != nil {
t.Fatal(err)
}
}
if err := FinishSybSyncRun(db, model.SybSyncRun{
RunID: "DONE", Status: model.SybSyncSucceeded, FinishedAt: "2026-08-09T01:01:00Z",
}); err != nil {
t.Fatal(err)
}
affected, err := InterruptRunningSybSyncRuns(db, "2026-08-09T02:00:00Z")
if err != nil || affected != 1 {
t.Fatalf("应该只中断一条 running 记录: affected=%d err=%v", affected, err)
}
rows, _ := ListSybSyncRuns(db, 10, 0)
statuses := map[string]model.SybSyncRunStatus{}
for _, row := range rows {
statuses[row.RunID] = row.Status
}
if statuses["RUNNING"] != model.SybSyncInterrupted || statuses["DONE"] != model.SybSyncSucceeded {
t.Fatalf("中断状态错误: %+v", statuses)
}
}
func TestListSybOrders_关键字筛选订单号和标题(t *testing.T) {
db := newSybTestDB(t)
mustUpsert := func(sybID, orderNo, title string) {