// 顺运宝会话缓存、同步进度和货运单明细的读写。 // // 改动前必读 admin/AGENTS.md:只有本文件(和 db.go)能写 SQL, // service/syb.go 和 handler/web/others.go 都不许拼 SQL。 package repository import ( "database/sql" "errors" "fmt" "strings" "cmautobuy/admin/model" ) // ---------- 会话缓存 ---------- // SaveSybSession 写入或更新顺运宝登录会话缓存(按用户名 upsert)。 // // `[必须]` 缓存写失败不能让已经登录的顺运宝会话失效——调用方拿到错误后 // 只应该记日志,不应该把内存里刚登录成功的会话也扔掉, // 见 docs/admin/08-顺运宝接口.md §8。这条约束在 service 层落实, // 这里只负责"写失败就如实返回错误"。 func SaveSybSession(q Execer, username, cookiesJSON, expiresAt string) error { if username == "" { return fmt.Errorf("username 不能为空") } now := model.NowISO() _, err := q.Exec(` INSERT INTO syb_session (username, cookies, expires_at, updated_at) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE cookies = VALUES(cookies), expires_at = VALUES(expires_at), updated_at = VALUES(updated_at)`, username, cookiesJSON, expiresAt, now) if err != nil { return fmt.Errorf("保存顺运宝会话缓存失败: %w", err) } return nil } // SybSessionCache 是缓存里的一条顺运宝会话。 type SybSessionCache struct { Username string Cookies string // JSON 数组,原样保留,交给 syb.Client 解析 ExpiresAt string } // GetSybSession 按用户名查缓存的会话,查不到返回 (nil, nil)。 // // `[必须]` 这里只负责取数据,**不判断是否过期**——过期时间的比较、 // 是否需要重新登录,是 service 层的业务判断(要用到"现在几点"这个 // 会变化的量,放这里测试起来还要控制时间,不如交给上层)。 func GetSybSession(q Execer, username string) (*SybSessionCache, error) { var c SybSessionCache err := q.QueryRow(` SELECT username, cookies, expires_at FROM syb_session WHERE username = ?`, username, ).Scan(&c.Username, &c.Cookies, &c.ExpiresAt) if errors.Is(err, sql.ErrNoRows) { return nil, nil } if err != nil { return nil, fmt.Errorf("查询顺运宝会话缓存失败: %w", err) } return &c, nil } // DeleteSybSession 清除某个用户名的会话缓存(会话确认失效后调用)。 func DeleteSybSession(q Execer, username string) error { if _, err := q.Exec(`DELETE FROM syb_session WHERE username = ?`, username); err != nil { return fmt.Errorf("清除顺运宝会话缓存失败: %w", err) } return nil } // ---------- 同步进度 ---------- // GetSybLastSyncedAt 查上次同步完成的时间(精确到秒的 ISO 字符串)。 // 从没同步过时返回 ("", false, nil)。 func GetSybLastSyncedAt(q Execer) (lastSyncedAt string, found bool, err error) { var s sql.NullString err = q.QueryRow(`SELECT last_synced_at FROM syb_sync_state WHERE id = 1`).Scan(&s) if errors.Is(err, sql.ErrNoRows) { return "", false, nil } if err != nil { return "", false, fmt.Errorf("查询顺运宝同步进度失败: %w", err) } if !s.Valid || s.String == "" { return "", false, nil } return s.String, true, nil } // SetSybLastSyncedAt 更新"上次同步到哪"。 // // `[必须]` 只应该在一次同步**全部成功**之后调用——调用方(service 层) // 负责这个时机;这里只负责写,不判断"是否该写"。中途失败不调用这个函数, // 见工单 #46「中途失败不更新 last_synced_at」。 func SetSybLastSyncedAt(q Execer, at string) error { now := model.NowISO() _, err := q.Exec(` INSERT INTO syb_sync_state (id, last_synced_at, updated_at) VALUES (1, ?, ?) ON DUPLICATE KEY UPDATE last_synced_at = VALUES(last_synced_at), updated_at = VALUES(updated_at)`, at, now) if err != nil { return fmt.Errorf("更新顺运宝同步进度失败: %w", err) } 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 写入或更新一条顺运宝货运单明细行。 // // `[必须]` DO UPDATE SET 里绝不允许出现 shopee_sku_id。它是规格匹配的 // 结果(人工确认或自动匹配产生),顺运宝那边根本没有这个值——写进去 // 就是写 NULL,把人工攒的匹配成果洗掉,而且不报错。这和 #38 里 // pdd_goods_url 不能被 Excel 导入覆盖是同一类问题,见工单 #46、 // docs/admin/08-顺运宝接口.md §6.2。 // // `[必须]` shopee_goods_id **可以**被覆盖——它就是顺运宝 // detail.productId,来自顺运宝,不是人工填的。 // // 返回 created 表示这一行是不是本次新插入的(供上层统计"新增/更新")。 func UpsertSybOrder(q Execer, o model.SybOrder) (created bool, err error) { if o.SybID == "" { return false, fmt.Errorf("syb_id 不能为空") } var exists int err = q.QueryRow(`SELECT 1 FROM syb_orders WHERE syb_id = ?`, o.SybID).Scan(&exists) switch { case errors.Is(err, sql.ErrNoRows): created = true case err != nil: return false, fmt.Errorf("查询顺运宝货运单明细 %s 失败: %w", o.SybID, err) } now := model.NowISO() _, err = q.Exec(` INSERT INTO syb_orders (syb_id, order_no, title, product_spec, shopee_goods_id, shopee_sku_id, quantity, price_twd_cent, image_url, syb_data, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE order_no = VALUES(order_no), title = VALUES(title), product_spec = VALUES(product_spec), shopee_goods_id = VALUES(shopee_goods_id), quantity = VALUES(quantity), price_twd_cent = VALUES(price_twd_cent), image_url = VALUES(image_url), syb_data = VALUES(syb_data), updated_at = VALUES(updated_at)`, // 注意:shopee_sku_id 只出现在 INSERT 的列清单里(新建行时写 o.ShopeeSKUID, // 同步永远传空字符串),完全不出现在 DO UPDATE SET 里——已存在的行 // 这一列不受本语句影响,见上面的函数注释。 o.SybID, o.OrderNo, o.Title, nullableText(o.ProductSpec), nullableText(o.ShopeeGoodsID), nullableText(o.ShopeeSKUID), o.Quantity, o.PriceTwdCent, nullableText(o.ImageURL), o.SybData, now, now) if err != nil { return false, fmt.Errorf("写入顺运宝货运单明细 %s 失败: %w", o.SybID, err) } return created, nil } // nullableText 把空字符串转成 SQL NULL,非空字符串原样写入。 // syb_orders 的这几列在建表语句里都允许 NULL,空字符串和 NULL // 在页面上显示效果一样,统一存 NULL 更符合"这个字段还没有值"的语义。 func nullableText(s string) any { if s == "" { return nil } return s } // SybOrderFilter 是货运单列表页支持的筛选条件。 type SybOrderFilter struct { Keyword string // 匹配订单号或商品标题 Stage string // 处理阶段;空字符串表示全部 } func sybOrderFilterClause(filter SybOrderFilter) (string, []any) { var clauses []string var args []any if kw := filter.Keyword; kw != "" { like := "%" + escapeLike(kw) + "%" clauses = append(clauses, `(so.order_no LIKE ? ESCAPE '!' OR so.title LIKE ? ESCAPE '!')`) args = append(args, like, like) } switch filter.Stage { case "missing_shopee": clauses = append(clauses, `sp.goods_id IS NULL`) case "sku_pending": clauses = append(clauses, `sp.goods_id IS NOT NULL AND so.shopee_sku_id IS NULL`) case "pdd_missing": clauses = append(clauses, `so.shopee_sku_id IS NOT NULL AND (sp.pdd_goods_id IS NULL OR sp.pdd_goods_id = '' OR pp.goods_id IS NULL)`) case "pdd_pending", "pdd_collecting", "pdd_failed", "pdd_collected": clauses = append(clauses, `so.shopee_sku_id IS NOT NULL AND pp.collect_status = ?`) args = append(args, strings.TrimPrefix(filter.Stage, "pdd_")) } if len(clauses) == 0 { return "", args } return ` WHERE ` + strings.Join(clauses, ` AND `), args } const sybOrderContextFrom = ` FROM syb_orders so LEFT JOIN shopee_products sp ON sp.goods_id = so.shopee_goods_id LEFT JOIN pdd_products pp ON pp.goods_id = sp.pdd_goods_id AND pp.deleted_at IS NULL LEFT JOIN sku_mappings sm ON sm.shopee_sku_id = so.shopee_sku_id AND sm.pdd_goods_id = sp.pdd_goods_id` // SybOrderContext 是顺运宝明细及其当前蝦皮/PDD 处理上下文。 // 处理阶段由 service 计算,Repository 只提供数据库事实。 type SybOrderContext struct { Order model.SybOrder ShopeeExists bool PddGoodsID string PddGoodsURL string PddCollectStatus string PddCollectMsg string PddSkusJSON string MappingOptionKey string MappingOptions string HasActiveTask bool } func scanSybOrderContext(s rowScanner) (SybOrderContext, error) { var c SybOrderContext var title, productSpec, shopeeGoodsID, shopeeSKUID, imageURL sql.NullString var priceCent sql.NullInt64 var shopeeExists int var pddGoodsID, pddGoodsURL, collectStatus, collectMsg, skusJSON sql.NullString var mappingKey, mappingOptions sql.NullString var hasActiveTask int err := s.Scan( &c.Order.SybID, &c.Order.OrderNo, &title, &productSpec, &shopeeGoodsID, &shopeeSKUID, &c.Order.Quantity, &priceCent, &imageURL, &c.Order.SybData, &c.Order.CreatedAt, &c.Order.UpdatedAt, &shopeeExists, &pddGoodsID, &pddGoodsURL, &collectStatus, &collectMsg, &skusJSON, &mappingKey, &mappingOptions, &hasActiveTask, ) c.Order.Title = title.String c.Order.ProductSpec = productSpec.String c.Order.ShopeeGoodsID = shopeeGoodsID.String c.Order.ShopeeSKUID = shopeeSKUID.String c.Order.PriceTwdCent = priceCent.Int64 c.Order.ImageURL = imageURL.String c.ShopeeExists = shopeeExists != 0 c.PddGoodsID = pddGoodsID.String c.PddGoodsURL = pddGoodsURL.String c.PddCollectStatus = collectStatus.String c.PddCollectMsg = collectMsg.String c.PddSkusJSON = skusJSON.String c.MappingOptionKey = mappingKey.String c.MappingOptions = mappingOptions.String c.HasActiveTask = hasActiveTask != 0 return c, err } // ListSybOrderContexts 分页读取采购处理工作台所需的关联事实。 func ListSybOrderContexts(q Execer, filter SybOrderFilter, limit, offset int) ([]SybOrderContext, error) { where, args := sybOrderFilterClause(filter) query := ` SELECT so.syb_id, so.order_no, so.title, so.product_spec, so.shopee_goods_id, so.shopee_sku_id, so.quantity, so.price_twd_cent, so.image_url, so.syb_data, so.created_at, so.updated_at, CASE WHEN sp.goods_id IS NULL THEN 0 ELSE 1 END, sp.pdd_goods_id, sp.pdd_goods_url, pp.collect_status, pp.collect_msg, pp.skus_json, sm.pdd_option_key, sm.pdd_options, EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase' AND t.syb_id = so.syb_id AND t.status IN ('pending', 'assigned', 'claimed'))` + sybOrderContextFrom + where + ` ORDER BY so.updated_at DESC, so.syb_id DESC` if limit >= 0 { query += ` LIMIT ? OFFSET ?` args = append(args, limit, offset) } rows, err := q.Query(query, args...) if err != nil { return nil, fmt.Errorf("查询顺运宝采购处理列表失败: %w", err) } defer rows.Close() list := make([]SybOrderContext, 0) for rows.Next() { row, err := scanSybOrderContext(rows) if err != nil { return nil, fmt.Errorf("读取顺运宝采购处理列表失败: %w", err) } list = append(list, row) } return list, rows.Err() } // GetSybOrderContext 按明细 ID 读取一条处理上下文。 func GetSybOrderContext(q Execer, sybID string) (*SybOrderContext, error) { row, err := scanSybOrderContext(q.QueryRow(` SELECT so.syb_id, so.order_no, so.title, so.product_spec, so.shopee_goods_id, so.shopee_sku_id, so.quantity, so.price_twd_cent, so.image_url, so.syb_data, so.created_at, so.updated_at, CASE WHEN sp.goods_id IS NULL THEN 0 ELSE 1 END, sp.pdd_goods_id, sp.pdd_goods_url, pp.collect_status, pp.collect_msg, pp.skus_json, sm.pdd_option_key, sm.pdd_options, EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase' AND t.syb_id = so.syb_id AND t.status IN ('pending', 'assigned', 'claimed'))`+ sybOrderContextFrom+` WHERE so.syb_id = ?`, sybID)) if errors.Is(err, sql.ErrNoRows) { return nil, nil } if err != nil { return nil, fmt.Errorf("查询顺运宝明细 %s 失败: %w", sybID, err) } return &row, nil } // SetSybShopeeSKUIfEmpty 只为尚未确认规格的明细写入自动识别结果。 func SetSybShopeeSKUIfEmpty(q Execer, sybID, skuID string) (bool, error) { result, err := q.Exec(` UPDATE syb_orders SET shopee_sku_id = ?, updated_at = ? WHERE syb_id = ? AND shopee_sku_id IS NULL AND EXISTS ( SELECT 1 FROM shopee_skus sk WHERE sk.sku_id = ? AND sk.goods_id = syb_orders.shopee_goods_id )`, skuID, model.NowISO(), sybID, skuID) if err != nil { return false, fmt.Errorf("自动确认顺运宝明细 %s 的蝦皮规格失败: %w", sybID, err) } n, err := result.RowsAffected() return n > 0, err } // SetSybShopeeSKU 在调用方完成商品归属校验后保存人工选择。 func SetSybShopeeSKU(q Execer, sybID, skuID string) error { result, err := q.Exec(`UPDATE syb_orders SET shopee_sku_id = ?, updated_at = ? WHERE syb_id = ?`, skuID, model.NowISO(), sybID) if err != nil { return fmt.Errorf("保存顺运宝明细 %s 的蝦皮规格失败: %w", sybID, err) } n, err := result.RowsAffected() if err != nil { return err } if n == 0 { return fmt.Errorf("顺运宝明细 %s 不存在", sybID) } return nil } // ListSybOrders 按筛选条件分页查货运单明细列表,按更新时间倒序。 func ListSybOrders(q Execer, filter SybOrderFilter, limit, offset int) ([]model.SybOrder, error) { where, args := sybOrderFilterClause(filter) sqlText := ` SELECT so.syb_id, so.order_no, so.title, so.product_spec, so.shopee_goods_id, so.shopee_sku_id, so.quantity, so.price_twd_cent, so.image_url, so.syb_data, so.created_at, so.updated_at` + sybOrderContextFrom + where + ` ORDER BY so.updated_at DESC, so.syb_id DESC LIMIT ? OFFSET ?` args = append(args, limit, offset) rows, err := q.Query(sqlText, args...) if err != nil { return nil, fmt.Errorf("查询顺运宝货运单列表失败: %w", err) } defer rows.Close() var list []model.SybOrder for rows.Next() { var o model.SybOrder var title, productSpec, shopeeGoodsID, shopeeSKUID, imageURL sql.NullString var priceCent sql.NullInt64 if err := rows.Scan( &o.SybID, &o.OrderNo, &title, &productSpec, &shopeeGoodsID, &shopeeSKUID, &o.Quantity, &priceCent, &imageURL, &o.SybData, &o.CreatedAt, &o.UpdatedAt, ); err != nil { return nil, fmt.Errorf("读取顺运宝货运单列表失败: %w", err) } o.Title = title.String o.ProductSpec = productSpec.String o.ShopeeGoodsID = shopeeGoodsID.String o.ShopeeSKUID = shopeeSKUID.String o.PriceTwdCent = priceCent.Int64 o.ImageURL = imageURL.String list = append(list, o) } return list, rows.Err() } // CountSybOrders 统计当前筛选条件下的货运单明细总数。 // // `[必须]` 用和 ListSybOrders **完全相同**的筛选条件——分页和底部统计 // 靠它,写成两份筛选条件迟早有一天会不一致(工单 #43 的教训)。 func CountSybOrders(q Execer, filter SybOrderFilter) (int, error) { where, args := sybOrderFilterClause(filter) sqlText := `SELECT COUNT(*)` + sybOrderContextFrom + where var n int if err := q.QueryRow(sqlText, args...).Scan(&n); err != nil { return 0, fmt.Errorf("统计顺运宝货运单数量失败: %w", err) } return n, nil } // CountSybOrdersTotal 统计全部货运单明细数量(不带筛选), // 供列表页判断"是否已经同步过任何数据"。 func CountSybOrdersTotal(q Execer) (int, error) { var n int if err := q.QueryRow(`SELECT COUNT(*) FROM syb_orders`).Scan(&n); err != nil { return 0, fmt.Errorf("统计顺运宝货运单数量失败: %w", err) } return n, nil }