813 lines
32 KiB
Go
813 lines
32 KiB
Go
// 顺运宝会话缓存、同步进度和货运单明细的读写。
|
|
//
|
|
// 改动前必读 admin/AGENTS.md:只有本文件(和 db.go)能写 SQL,
|
|
// service/syb.go 和 handler/web/others.go 都不许拼 SQL。
|
|
package repository
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"database/sql"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
|
|
"cmautobuy/admin/model"
|
|
"cmautobuy/admin/spec"
|
|
)
|
|
|
|
// ---------- 会话缓存 ----------
|
|
|
|
// 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 = ?, 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.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)
|
|
}
|
|
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.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
|
|
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, 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.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)
|
|
}
|
|
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 写入或更新一条顺运宝货运单明细行。
|
|
//
|
|
// `[必须]` shopee_sku_id 已退出采购主链路,INSERT 和 UPDATE 都不再写它。
|
|
// 旧列只为只追加迁移兼容保留。
|
|
//
|
|
// `[必须]` 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 不能为空")
|
|
}
|
|
if strings.TrimSpace(o.ShopID) == "" && strings.TrimSpace(o.ShopName) != "" {
|
|
o.ShopID, err = FindShopIDByName(q, o.ShopName)
|
|
if err != nil {
|
|
return false, fmt.Errorf("解析顺运宝明细 %s 的店铺失败: %w", o.SybID, err)
|
|
}
|
|
}
|
|
var specKey any
|
|
if strings.TrimSpace(o.ProductSpec) != "" {
|
|
key, keyErr := spec.SpecKey(o.ProductSpec)
|
|
if keyErr != nil {
|
|
return false, fmt.Errorf("顺运宝明细 %s 的规格不能生成身份键: %w", o.SybID, keyErr)
|
|
}
|
|
specKey = key
|
|
}
|
|
|
|
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()
|
|
if strings.TrimSpace(o.ShopeeGoodsID) != "" {
|
|
if err := upsertSybShopeeProduct(q, o.ShopeeGoodsID, o.ShopID, o.Title, o.ShopName, o.ImageURL, now); err != nil {
|
|
return false, fmt.Errorf("为顺运宝明细 %s 补建蝦皮商品骨架失败: %w", o.SybID, err)
|
|
}
|
|
}
|
|
_, err = q.Exec(`
|
|
INSERT INTO syb_orders
|
|
(syb_id, order_no, shop_id, shop_name, title, product_spec, spec_key, shopee_goods_id,
|
|
quantity, price_twd_cent, image_url, syb_data, created_at, updated_at)
|
|
VALUES (?, ?, NULLIF(?,''), ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON DUPLICATE KEY UPDATE
|
|
order_no = VALUES(order_no),
|
|
shop_id = VALUES(shop_id),
|
|
shop_name = VALUES(shop_name),
|
|
title = VALUES(title),
|
|
product_spec = VALUES(product_spec),
|
|
spec_key = VALUES(spec_key),
|
|
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)`,
|
|
o.SybID, o.OrderNo, strings.TrimSpace(o.ShopID), nullableText(o.ShopName), o.Title, nullableText(o.ProductSpec), specKey, nullableText(o.ShopeeGoodsID),
|
|
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
|
|
}
|
|
|
|
// upsertSybShopeeProduct 用顺运宝观测补建蝦皮商品骨架。
|
|
//
|
|
// 店铺和图片是字段级低优先级数据:只能补空值,或更新原本同样来自 syb 的值;
|
|
// 人工字段和商品目录等权威来源永远不被顺运宝覆盖。空值也不能清除已有内容。
|
|
func upsertSybShopeeProduct(q Execer, goodsID, shopID, title, shopName, imageURL, observedAt string) error {
|
|
_, err := q.Exec(`
|
|
INSERT INTO shopee_products
|
|
(goods_id,shop_id,title,image_url,shopee_shop_name,
|
|
image_source,image_observed_at,image_is_manual,
|
|
shop_name_source,shop_name_observed_at,shop_name_is_manual,
|
|
source,source_observed_at,created_at,updated_at)
|
|
VALUES (?, NULLIF(TRIM(?),''), ?, NULLIF(TRIM(?),''), NULLIF(TRIM(?),''),
|
|
CASE WHEN TRIM(?)='' THEN NULL ELSE 'syb' END,
|
|
CASE WHEN TRIM(?)='' THEN NULL ELSE ? END, 0,
|
|
CASE WHEN TRIM(?)='' THEN NULL ELSE 'syb' END,
|
|
CASE WHEN TRIM(?)='' THEN NULL ELSE ? END, 0,
|
|
'syb', ?, ?, ?)
|
|
ON DUPLICATE KEY UPDATE
|
|
shop_id = CASE
|
|
WHEN VALUES(shop_id) IS NOT NULL AND (shop_id IS NULL OR shop_name_source='syb') THEN VALUES(shop_id)
|
|
ELSE shop_id
|
|
END,
|
|
title = CASE
|
|
WHEN source='syb' AND TRIM(VALUES(title))<>'' THEN VALUES(title)
|
|
ELSE title
|
|
END,
|
|
source_observed_at = CASE
|
|
WHEN source='syb' AND TRIM(VALUES(title))<>'' THEN VALUES(source_observed_at)
|
|
ELSE source_observed_at
|
|
END,
|
|
image_url = CASE
|
|
WHEN image_is_manual=0 AND TRIM(VALUES(image_url))<>''
|
|
AND (image_url IS NULL OR TRIM(image_url)='' OR image_source='syb') THEN VALUES(image_url)
|
|
ELSE image_url
|
|
END,
|
|
image_observed_at = CASE
|
|
WHEN image_is_manual=0 AND TRIM(VALUES(image_url))<>''
|
|
AND (image_url IS NULL OR TRIM(image_url)='' OR image_source='syb') THEN VALUES(image_observed_at)
|
|
ELSE image_observed_at
|
|
END,
|
|
image_source = CASE
|
|
WHEN image_is_manual=0 AND TRIM(VALUES(image_url))<>''
|
|
AND (image_url IS NULL OR TRIM(image_url)='' OR image_source='syb') THEN 'syb'
|
|
ELSE image_source
|
|
END,
|
|
shopee_shop_name = CASE
|
|
WHEN shop_name_is_manual=0 AND TRIM(VALUES(shopee_shop_name))<>''
|
|
AND (shopee_shop_name IS NULL OR TRIM(shopee_shop_name)='' OR shop_name_source='syb') THEN VALUES(shopee_shop_name)
|
|
ELSE shopee_shop_name
|
|
END,
|
|
shop_name_observed_at = CASE
|
|
WHEN shop_name_is_manual=0 AND TRIM(VALUES(shopee_shop_name))<>''
|
|
AND (shopee_shop_name IS NULL OR TRIM(shopee_shop_name)='' OR shop_name_source='syb') THEN VALUES(shop_name_observed_at)
|
|
ELSE shop_name_observed_at
|
|
END,
|
|
shop_name_source = CASE
|
|
WHEN shop_name_is_manual=0 AND TRIM(VALUES(shopee_shop_name))<>''
|
|
AND (shopee_shop_name IS NULL OR TRIM(shopee_shop_name)='' OR shop_name_source='syb') THEN 'syb'
|
|
ELSE shop_name_source
|
|
END,
|
|
updated_at = CASE
|
|
WHEN source='syb'
|
|
OR (image_source='syb' AND TRIM(VALUES(image_url))<>'')
|
|
OR (shop_name_source='syb' AND TRIM(VALUES(shopee_shop_name))<>'') THEN VALUES(updated_at)
|
|
ELSE updated_at
|
|
END`,
|
|
goodsID, shopID, title, imageURL, shopName,
|
|
imageURL, imageURL, observedAt,
|
|
shopName, shopName, observedAt,
|
|
observedAt, observedAt, observedAt)
|
|
return err
|
|
}
|
|
|
|
// backfillSybProductMetadata 把历史货运单中每个商品最新的非空店铺、图片补入商品主表。
|
|
// 重放是安全的:实际覆盖规则仍由 upsertSybShopeeProduct 统一执行。
|
|
func backfillSybProductMetadata(db *sql.DB) error {
|
|
type metadata struct{ shopName, imageURL, observedAt string }
|
|
items := map[string]metadata{}
|
|
readLatest := func(column string, assign func(*metadata, string)) error {
|
|
query := fmt.Sprintf(`SELECT shopee_goods_id, value_text, updated_at FROM (
|
|
SELECT shopee_goods_id, %s AS value_text, updated_at,
|
|
ROW_NUMBER() OVER (PARTITION BY shopee_goods_id ORDER BY updated_at DESC, syb_id DESC) AS row_no
|
|
FROM syb_orders
|
|
WHERE TRIM(COALESCE(shopee_goods_id,''))<>'' AND TRIM(COALESCE(%s,''))<>''
|
|
) ranked WHERE row_no=1`, column, column)
|
|
rows, err := db.Query(query)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer rows.Close()
|
|
for rows.Next() {
|
|
var goodsID, value, observedAt string
|
|
if err := rows.Scan(&goodsID, &value, &observedAt); err != nil {
|
|
return err
|
|
}
|
|
item := items[goodsID]
|
|
assign(&item, value)
|
|
if observedAt > item.observedAt {
|
|
item.observedAt = observedAt
|
|
}
|
|
items[goodsID] = item
|
|
}
|
|
return rows.Err()
|
|
}
|
|
if err := readLatest("shop_name", func(item *metadata, value string) { item.shopName = value }); err != nil {
|
|
return fmt.Errorf("读取历史顺运宝店铺失败: %w", err)
|
|
}
|
|
if err := readLatest("image_url", func(item *metadata, value string) { item.imageURL = value }); err != nil {
|
|
return fmt.Errorf("读取历史顺运宝图片失败: %w", err)
|
|
}
|
|
for goodsID, item := range items {
|
|
shopID, err := FindShopIDByName(db, item.shopName)
|
|
if err != nil {
|
|
return fmt.Errorf("解析历史顺运宝店铺失败: %w", err)
|
|
}
|
|
if err := upsertSybShopeeProduct(db, goodsID, shopID, "", item.shopName, item.imageURL, item.observedAt); err != nil {
|
|
return fmt.Errorf("回填蝦皮商品 %s 的顺运宝元数据失败: %w", goodsID, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// UpsertSybShopeeSKU 把格式明确的 SYB 规格作为低优先级蝦皮 SKU 观测写入。
|
|
// 返回 parsed=false 时不写库,避免把猜测结果污染商品主数据。
|
|
func UpsertSybShopeeSKU(q Execer, goodsID, specRaw, observedAt string) (parsed bool, err error) {
|
|
parsedSpec, ok := spec.ParseShopeeSpec(specRaw)
|
|
if !ok || strings.TrimSpace(goodsID) == "" {
|
|
return false, nil
|
|
}
|
|
specKey, err := spec.SpecKey(specRaw)
|
|
if err != nil {
|
|
return false, nil
|
|
}
|
|
sum := sha256.Sum256([]byte(goodsID + "\x00" + specKey))
|
|
recordID := "syb:" + hex.EncodeToString(sum[:])
|
|
now := model.NowISO()
|
|
_, err = UpsertCatalogShopeeSKU(q, CatalogShopeeSKUInput{
|
|
RecordID: recordID, GoodsID: goodsID, SpecRaw: specRaw, SpecKey: specKey,
|
|
Color: parsedSpec.Color, Size: parsedSpec.Size, Advice: parsedSpec.Advice,
|
|
ParseOK: true, Source: "syb", ObservedAt: observedAt, Now: now, UpdatePolicy: "fill_missing",
|
|
})
|
|
return true, err
|
|
}
|
|
|
|
func backfillSybShopeeSKUs(db *sql.DB) error {
|
|
rows, err := db.Query(`SELECT shopee_goods_id,product_spec,updated_at FROM (
|
|
SELECT shopee_goods_id,product_spec,updated_at,
|
|
ROW_NUMBER() OVER (PARTITION BY shopee_goods_id,product_spec ORDER BY updated_at DESC,syb_id DESC) row_no
|
|
FROM syb_orders
|
|
WHERE TRIM(COALESCE(shopee_goods_id,''))<>'' AND TRIM(COALESCE(product_spec,''))<>''
|
|
) ranked WHERE row_no=1`)
|
|
if err != nil {
|
|
return fmt.Errorf("读取历史顺运宝规格失败: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
type observation struct{ goodsID, raw, observedAt string }
|
|
var observations []observation
|
|
for rows.Next() {
|
|
var item observation
|
|
if err := rows.Scan(&item.goodsID, &item.raw, &item.observedAt); err != nil {
|
|
return err
|
|
}
|
|
observations = append(observations, item)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return err
|
|
}
|
|
for _, item := range observations {
|
|
if _, err := UpsertSybShopeeSKU(db, item.goodsID, item.raw, item.observedAt); err != nil {
|
|
return fmt.Errorf("回填蝦皮商品 %s 的顺运宝规格失败: %w", item.goodsID, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SybSpecObservation 是顺运宝对某个蝦皮商品规格的历史观测汇总。
|
|
// 它不是正式蝦皮 SKU,不得写回 shopee_skus。
|
|
type SybSpecObservation struct {
|
|
SpecKey string
|
|
SpecRaw string
|
|
LatestPriceCent int64
|
|
LatestImageURL string
|
|
OrderCount int
|
|
TotalQuantity int
|
|
LastObservedAt string
|
|
}
|
|
|
|
// ListSybSpecObservations 按规格身份汇总货运单明细;同一 syb_id 重复同步仍只有一行。
|
|
func ListSybSpecObservations(q Execer, goodsID string) ([]SybSpecObservation, error) {
|
|
rows, err := q.Query(`WITH ranked AS (
|
|
SELECT spec_key,product_spec,price_twd_cent,image_url,quantity,created_at,syb_id,
|
|
ROW_NUMBER() OVER(PARTITION BY spec_key ORDER BY created_at DESC,syb_id DESC) AS row_num,
|
|
COUNT(*) OVER(PARTITION BY spec_key) AS order_count,
|
|
SUM(quantity) OVER(PARTITION BY spec_key) AS total_quantity,
|
|
MAX(created_at) OVER(PARTITION BY spec_key) AS last_observed_at
|
|
FROM syb_orders
|
|
WHERE shopee_goods_id=? AND spec_key IS NOT NULL AND spec_key<>''
|
|
)
|
|
SELECT spec_key,product_spec,price_twd_cent,image_url,order_count,total_quantity,last_observed_at
|
|
FROM ranked WHERE row_num=1 ORDER BY last_observed_at DESC,spec_key`, goodsID)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("查询蝦皮商品 %s 的顺运宝观测规格失败: %w", goodsID, err)
|
|
}
|
|
defer rows.Close()
|
|
var list []SybSpecObservation
|
|
for rows.Next() {
|
|
var item SybSpecObservation
|
|
var raw, image sql.NullString
|
|
if err := rows.Scan(&item.SpecKey, &raw, &item.LatestPriceCent, &image, &item.OrderCount, &item.TotalQuantity, &item.LastObservedAt); err != nil {
|
|
return nil, fmt.Errorf("读取顺运宝观测规格失败: %w", err)
|
|
}
|
|
item.SpecRaw = raw.String
|
|
item.LatestImageURL = image.String
|
|
list = append(list, item)
|
|
}
|
|
return list, rows.Err()
|
|
}
|
|
|
|
// nullableText 把空字符串转成 SQL NULL,非空字符串原样写入。
|
|
// syb_orders 的这几列在建表语句里都允许 NULL,空字符串和 NULL
|
|
// 在页面上显示效果一样,统一存 NULL 更符合"这个字段还没有值"的语义。
|
|
func nullableText(s string) any {
|
|
if s == "" {
|
|
return nil
|
|
}
|
|
return s
|
|
}
|
|
|
|
// SybOrderFilter 是货运单列表页支持的筛选条件。
|
|
type SybOrderFilter struct {
|
|
OrderNos []string // 订单号精确匹配;多个值之间是 OR,空切片表示全部
|
|
Shop string // 店铺名模糊匹配;空字符串表示全部
|
|
Stage string // 处理阶段;空字符串表示全部
|
|
}
|
|
|
|
func sybOrderFilterClause(filter SybOrderFilter) (string, []any) {
|
|
var clauses []string
|
|
var args []any
|
|
if len(filter.OrderNos) > 0 {
|
|
placeholders := make([]string, len(filter.OrderNos))
|
|
for i, orderNo := range filter.OrderNos {
|
|
placeholders[i] = "?"
|
|
args = append(args, orderNo)
|
|
}
|
|
clauses = append(clauses, `so.order_no IN (`+strings.Join(placeholders, ",")+`)`)
|
|
}
|
|
if shop := strings.TrimSpace(filter.Shop); shop != "" {
|
|
clauses = append(clauses, `so.shop_name LIKE ? ESCAPE '!'`)
|
|
args = append(args, "%"+escapeLike(shop)+"%")
|
|
}
|
|
switch filter.Stage {
|
|
case "spec_missing":
|
|
clauses = append(clauses, `so.spec_key IS NULL`)
|
|
case "pdd_missing":
|
|
clauses = append(clauses, `so.spec_key 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.spec_key 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 AND sp.deleted_at IS NULL
|
|
LEFT JOIN pdd_products pp ON pp.goods_id = sp.pdd_goods_id AND pp.deleted_at IS NULL
|
|
LEFT JOIN spec_mappings sm
|
|
ON sm.shopee_goods_id = so.shopee_goods_id
|
|
AND sm.spec_key = so.spec_key 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
|
|
PddUpdatedAt string
|
|
MappingOptionKey string
|
|
MappingOptions string
|
|
MappingSource string
|
|
MappingProviderID string
|
|
MappingModel string
|
|
MappingConfidenceBPS int
|
|
MappingReason string
|
|
MappingSourceVersion string
|
|
MappingContextVersion string
|
|
HasSucceededPurchaseTask bool
|
|
HasManualReviewPurchaseTask bool
|
|
HasActiveTask bool
|
|
HasActiveCollectTask bool
|
|
}
|
|
|
|
func scanSybOrderContext(s rowScanner) (SybOrderContext, error) {
|
|
var c SybOrderContext
|
|
var shopName, title, productSpec, specKey, shopeeGoodsID, imageURL sql.NullString
|
|
var priceCent sql.NullInt64
|
|
var shopeeExists int
|
|
var pddGoodsID, pddGoodsURL, collectStatus, collectMsg, skusJSON, pddUpdatedAt sql.NullString
|
|
var mappingKey, mappingOptions, mappingSource, mappingProviderID, mappingModel sql.NullString
|
|
var mappingReason, mappingSourceVersion, mappingContextVersion sql.NullString
|
|
var mappingConfidence sql.NullInt64
|
|
var hasSucceededTask, hasManualReviewTask, hasActiveTask, hasActiveCollectTask int
|
|
err := s.Scan(
|
|
&c.Order.SybID, &c.Order.OrderNo, &shopName, &title, &productSpec, &specKey, &shopeeGoodsID,
|
|
&c.Order.Quantity, &priceCent, &imageURL, &c.Order.SybData,
|
|
&c.Order.CreatedAt, &c.Order.UpdatedAt, &shopeeExists,
|
|
&pddGoodsID, &pddGoodsURL, &collectStatus, &collectMsg, &skusJSON, &pddUpdatedAt,
|
|
&mappingKey, &mappingOptions, &mappingSource, &mappingProviderID, &mappingModel,
|
|
&mappingConfidence, &mappingReason, &mappingSourceVersion, &mappingContextVersion,
|
|
&hasSucceededTask, &hasManualReviewTask,
|
|
&hasActiveTask, &hasActiveCollectTask,
|
|
)
|
|
c.Order.ShopName = shopName.String
|
|
c.Order.Title = title.String
|
|
c.Order.ProductSpec = productSpec.String
|
|
c.Order.SpecKey = specKey.String
|
|
c.Order.ShopeeGoodsID = shopeeGoodsID.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.PddUpdatedAt = pddUpdatedAt.String
|
|
c.MappingOptionKey = mappingKey.String
|
|
c.MappingOptions = mappingOptions.String
|
|
c.MappingSource = mappingSource.String
|
|
c.MappingProviderID = mappingProviderID.String
|
|
c.MappingModel = mappingModel.String
|
|
c.MappingConfidenceBPS = int(mappingConfidence.Int64)
|
|
c.MappingReason = mappingReason.String
|
|
c.MappingSourceVersion = mappingSourceVersion.String
|
|
c.MappingContextVersion = mappingContextVersion.String
|
|
c.HasSucceededPurchaseTask = hasSucceededTask != 0
|
|
c.HasManualReviewPurchaseTask = hasManualReviewTask != 0
|
|
c.HasActiveTask = hasActiveTask != 0
|
|
c.HasActiveCollectTask = hasActiveCollectTask != 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.shop_name, so.title, so.product_spec, so.spec_key,
|
|
so.shopee_goods_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, pp.updated_at, sm.pdd_option_key, sm.pdd_options,
|
|
sm.source,sm.source_provider_id,sm.source_model,sm.confidence_bps,
|
|
sm.source_reason,sm.source_version,sm.context_version,
|
|
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
|
|
AND t.syb_id = so.syb_id AND t.status = 'succeeded'),
|
|
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
|
|
AND t.syb_id = so.syb_id AND t.status = 'manual_review'),
|
|
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')),
|
|
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'collect'
|
|
AND t.pdd_goods_id = pp.goods_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.shop_name, so.title, so.product_spec, so.spec_key,
|
|
so.shopee_goods_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, pp.updated_at, sm.pdd_option_key, sm.pdd_options,
|
|
sm.source,sm.source_provider_id,sm.source_model,sm.confidence_bps,
|
|
sm.source_reason,sm.source_version,sm.context_version,
|
|
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
|
|
AND t.syb_id = so.syb_id AND t.status = 'succeeded'),
|
|
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'purchase'
|
|
AND t.syb_id = so.syb_id AND t.status = 'manual_review'),
|
|
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')),
|
|
EXISTS(SELECT 1 FROM tasks t WHERE t.task_type = 'collect'
|
|
AND t.pdd_goods_id = pp.goods_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
|
|
}
|
|
|
|
// 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.shop_name, so.title, so.product_spec, so.spec_key, so.shopee_goods_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 shopName, title, productSpec, specKey, shopeeGoodsID, imageURL sql.NullString
|
|
var priceCent sql.NullInt64
|
|
if err := rows.Scan(
|
|
&o.SybID, &o.OrderNo, &shopName, &title, &productSpec, &specKey, &shopeeGoodsID,
|
|
&o.Quantity, &priceCent, &imageURL, &o.SybData, &o.CreatedAt, &o.UpdatedAt,
|
|
); err != nil {
|
|
return nil, fmt.Errorf("读取顺运宝货运单列表失败: %w", err)
|
|
}
|
|
o.ShopName = shopName.String
|
|
o.Title = title.String
|
|
o.ProductSpec = productSpec.String
|
|
o.SpecKey = specKey.String
|
|
o.ShopeeGoodsID = shopeeGoodsID.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
|
|
}
|
|
|
|
// ListMatchedSybOrderNos 返回当前可由 SQL 表达的组合筛选实际命中的订单号。
|
|
// 页面用它区分“输入了多少订单”和“命中了多少商品明细”;查询继续复用
|
|
// idx_syb_orders_order,所有订单号都通过占位符传入。
|
|
func ListMatchedSybOrderNos(q Execer, filter SybOrderFilter) ([]string, error) {
|
|
if len(filter.OrderNos) == 0 {
|
|
return nil, nil
|
|
}
|
|
where, args := sybOrderFilterClause(filter)
|
|
rows, err := q.Query(`SELECT DISTINCT so.order_no`+sybOrderContextFrom+where+` ORDER BY so.order_no`, args...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("查询顺运宝命中订单号失败: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
matched := make([]string, 0, len(filter.OrderNos))
|
|
for rows.Next() {
|
|
var orderNo string
|
|
if err := rows.Scan(&orderNo); err != nil {
|
|
return nil, fmt.Errorf("读取顺运宝命中订单号失败: %w", err)
|
|
}
|
|
matched = append(matched, orderNo)
|
|
}
|
|
return matched, rows.Err()
|
|
}
|
|
|
|
// 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
|
|
}
|