Files
cmautobuy/admin/service/syb.go
T

1583 lines
58 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 顺运宝货运单同步的编排逻辑:算日期范围、翻页拉列表和明细、
// 字段映射、落库统计、汇总报告。
//
// 改动前必读 admin/AGENTS.md:本层不认识 *gin.Context,也不拼 SQL——
// 那些分别在 handler/web/others.go 和 repository/syb.go。
//
// 接口契约见 docs/admin/08-顺运宝接口.md,工单见 #46。三条最容易出事的规则:
// 1. shopee_sku_id 是历史兼容列,同步不读写;规格身份由 product_spec/spec_key 表达。
// 2. 增量必须从"上次同步日期当天"重新拉,不是第二天,见 syncDateRange。
// 3. 中途失败不更新 last_synced_at,见 RunSybSync 的最后一步。
package service
import (
"context"
"crypto/sha256"
"database/sql"
"encoding/json"
"errors"
"fmt"
"math"
"sort"
"strconv"
"strings"
"sync"
"time"
"unicode"
"unicode/utf8"
"cmautobuy/admin/config"
"cmautobuy/admin/model"
"cmautobuy/admin/repository"
"cmautobuy/admin/syb"
)
func mappingContextVersion(c repository.SybOrderContext) string {
payload, _ := json.Marshal([]string{c.Order.SybID, c.Order.SpecKey, c.Order.UpdatedAt,
c.Order.ProductSpec, strconv.Itoa(c.Order.Quantity), c.PddGoodsID, c.PddUpdatedAt,
c.PddSkusJSON, SpecMatchRulesVersion})
return fmt.Sprintf("%x", sha256.Sum256(payload))
}
const (
dateLayout = "2006-01-02"
maxSpecifiedSyncDays = 31
defaultSybMaxMatches = 10000
sybTodayListMaxAttempts = 3
)
// SybSyncOptions 表示一次同步的日期选择。页面必须明确传入两端日期;两端都空
// 仅保留给旧版内部调用兼容。
//
// `[必须]` 不允许只填一端。把“半个范围”悄悄退化成自动增量,会让操作员
// 以为补拉了指定日期,实际却跑了另一段数据。
type SybSyncOptions struct {
From string
To string
}
// IsSpecified 判断是否明确给出了日期范围。
func (o SybSyncOptions) IsSpecified() bool {
return strings.TrimSpace(o.From) != "" || strings.TrimSpace(o.To) != ""
}
// NewSybSyncOptions 清理并校验浏览器提交的日期范围。
// 日期按顺运宝服务端 UTC+8 解释,起止两天都包含。
func NewSybSyncOptions(fromRaw, toRaw string, now time.Time) (SybSyncOptions, error) {
from := strings.TrimSpace(fromRaw)
to := strings.TrimSpace(toRaw)
if from == "" && to == "" {
return SybSyncOptions{}, nil
}
if from == "" || to == "" {
return SybSyncOptions{}, fmt.Errorf("开始日期和结束日期必须同时填写")
}
fromDate, err := time.Parse(dateLayout, from)
if err != nil || fromDate.Format(dateLayout) != from {
return SybSyncOptions{}, fmt.Errorf("开始日期格式不正确,请使用 YYYY-MM-DD")
}
toDate, err := time.Parse(dateLayout, to)
if err != nil || toDate.Format(dateLayout) != to {
return SybSyncOptions{}, fmt.Errorf("结束日期格式不正确,请使用 YYYY-MM-DD")
}
if fromDate.After(toDate) {
return SybSyncOptions{}, fmt.Errorf("开始日期不能晚于结束日期")
}
if to > dateOf(now) {
return SybSyncOptions{}, fmt.Errorf("结束日期不能晚于顺运宝服务端今天(%s)", dateOf(now))
}
days := int(toDate.Sub(fromDate).Hours()/24) + 1
if days > maxSpecifiedSyncDays {
return SybSyncOptions{}, fmt.Errorf("同步日期范围最多选择 %d 天,当前范围为 %d 天,请分段同步",
maxSpecifiedSyncDays, days)
}
return SybSyncOptions{From: from, To: to}, nil
}
// splitDateRange 把闭区间拆成逐日查询段。顺运宝只支持日期粒度,按天拉取
// 能把 offset 分页期间新增数据造成的位移限制在当天。
func splitDateRange(from, to string) ([]string, error) {
fromDate, err := time.Parse(dateLayout, from)
if err != nil || fromDate.Format(dateLayout) != from {
return nil, fmt.Errorf("开始日期格式不正确,请使用 YYYY-MM-DD")
}
toDate, err := time.Parse(dateLayout, to)
if err != nil || toDate.Format(dateLayout) != to {
return nil, fmt.Errorf("结束日期格式不正确,请使用 YYYY-MM-DD")
}
if fromDate.After(toDate) {
return nil, fmt.Errorf("开始日期不能晚于结束日期")
}
dates := make([]string, 0, int(toDate.Sub(fromDate).Hours()/24)+1)
for day := fromDate; !day.After(toDate); day = day.AddDate(0, 0, 1) {
dates = append(dates, day.Format(dateLayout))
}
return dates, nil
}
// sybLocation 是顺运宝服务端的时区。
//
// `[必须]` 日期范围筛的是服务端的 created,而它是 UTC+8——实测
// raw_data/shunyunbaoerp_stock_query.har:抓包于 2026-07-28T03:31:45Z
// (= 11:31 UTC+8),同一响应里 created 是 "2026-07-28 10:37:59"。
// 若 created 是 UTC,换算成 UTC+8 就是 18:37,比抓包时刻晚 7 小时,
// 订单创建于未来,不成立;作为 UTC+8 讲得通(比抓包早 54 分钟)。
//
// 用 UTC 算日期会在本地(UTC+8)00:00–08:00 这段时间把"今天"算成昨天,
// 当天早晨创建的单这一轮拉不到(下一轮的 from 仍是上次同步日,
// 范围会覆盖回来,不会永久丢,但操作员当场会以为同步坏了)。
// 见 docs/admin/08-顺运宝接口.md §5.3。
//
// `[必须]` 用 time.FixedZone 写死,不用 time.LoadLocation("Asia/Shanghai")——
// 那要读系统 tzdata,Windows 默认没有,打包成 exe 后会失败。
var sybLocation = time.FixedZone("UTC+8", 8*60*60)
// dateOf 把一个时刻转成顺运宝服务端时区(UTC+8)下的 YYYY-MM-DD。
//
// `[必须]` last_synced_at 存的仍是 UTC ISO(和全库其它时间戳一致,
// model.NowISO 的约定不动),只在这里换算成 UTC+8 取日期。
func dateOf(t time.Time) string {
return t.In(sybLocation).Format(dateLayout)
}
// SybToday 返回顺运宝服务端时区下的今天,供页面 date 输入的 max/default
// 使用。页面和同步校验共用这一处,避免本机时区不同导致显示和校验差一天。
func SybToday(now time.Time) string {
return dateOf(now)
}
// SybSyncDefaults 是页面首次打开时预填的同步范围和必要提示。
type SybSyncDefaults struct {
From string
To string
Warning string
}
// DefaultSybSyncRange 固定预填昨天到今天。覆盖游标仍由同步服务维护,但不再
// 改写采购员眼前的日期选择;需要补历史缺口时由采购员明确选择日期范围。
func DefaultSybSyncRange(_ *sql.DB, now time.Time) (SybSyncDefaults, error) {
today, err := time.Parse(dateLayout, dateOf(now))
if err != nil {
return SybSyncDefaults{}, fmt.Errorf("计算顺运宝今天日期失败: %w", err)
}
return SybSyncDefaults{
From: today.AddDate(0, 0, -1).Format(dateLayout),
To: today.Format(dateLayout),
}, nil
}
// cursorISOForDate 把“已连续覆盖到哪一天”保存成 UTC ISO。下一次仍从该日
// 重拉,宁可重复当天,也不能漏掉该日稍晚创建的货运单。
func cursorISOForDate(date string) (string, error) {
t, err := time.ParseInLocation(dateLayout, date, sybLocation)
if err != nil || t.Format(dateLayout) != date {
return "", fmt.Errorf("同步游标日期格式不正确: %q", date)
}
return t.UTC().Format(model.TimeLayout), nil
}
// syncDateRange 算出这次同步该拉哪个日期范围。
//
// `[必须]` 增量必须从"上次同步日期当天"重新拉,不是第二天——created
// 筛选粒度是日期,last_synced_at 精确到秒,从第二天拉会漏掉当天晚些
// 时候创建的单,且不会报错。宁可重复拉(靠 upsert 幂等)也不能漏,
// 见工单 #46。
//
// `[必须]` 首次同步(lastSyncedAt 为空)用 syncFrom;结束日期用 now
// 对应的日期,不用未来日期。
func syncDateRange(lastSyncedAt string, syncFrom string, now time.Time) (from, to string, err error) {
to = dateOf(now)
if lastSyncedAt == "" {
from = strings.TrimSpace(syncFrom)
if from == "" {
return "", "", fmt.Errorf("从未同步过,且 config.yaml 里没有配置 syb.sync_from,无法确定起始日期")
}
return from, to, nil
}
t, ok := model.ParseISO(lastSyncedAt)
if !ok {
return "", "", fmt.Errorf("上次同步时间 %q 解析失败", lastSyncedAt)
}
return dateOf(t), to, nil
}
// ---------- 同步报告 ----------
// SkipNote 是一条跳过或失败的说明。
//
// `[必须]` 有跳过或失败时要把它们列出来,不能只给个数字,见工单 #46
// 「报告要说清楚」。
type SkipNote struct {
SybID string
Reason string
}
// SyncReport 是一次同步的结果,供状态条显示。
type SyncReport struct {
From, To string
Specified bool // true 表示操作员发起的指定日期补同步
StockCount int // 拉到的货运单数
AcceptedCount int // 店铺准入后接受的货运单数
ShopSkipped int // 店铺不在允许列表或为空而跳过的货运单数
ShopFilterHash string // 本次固定店铺快照的 SHA-256,不保存敏感凭据
DetailCount int // 落库的商品明细行数(不含跳过的)
Created int
Updated int
SkippedZero int // quantity <= 0 被跳过的条数
Notes []SkipNote
Err error
StartedAt time.Time
FinishedAt time.Time
CursorAdvanced bool
}
const sybSyncHistoryPageSize = 10
// SybSyncRunView 是同步记录弹窗的一行,时间和状态已经转成采购员可读文本。
type SybSyncRunView struct {
RunID string
Username string
DateRange string
StatusText string
StatusClass string
Summary string
ErrorMessage string
CursorText string
StartedAt string
FinishedAt string
}
// SybSyncHistoryResult 是同步记录弹窗的分页结果。
type SybSyncHistoryResult struct {
Rows []SybSyncRunView
Page int
TotalPages int
Total int
}
// CreateSybSyncRun 在后台任务启动前创建可审计记录。
func CreateSybSyncRun(db *sql.DB, actor *model.User, options SybSyncOptions, now time.Time) (string, error) {
if actor == nil || actor.UserID == "" {
return "", fmt.Errorf("无法确认当前操作账号,同步没有启动")
}
if !options.IsSpecified() {
return "", fmt.Errorf("同步必须明确填写开始日期和结束日期")
}
runID, err := randomID("SYB-", 16)
if err != nil {
return "", fmt.Errorf("生成同步记录编号失败: %w", err)
}
run := model.SybSyncRun{
RunID: runID, UserID: actor.UserID, DateFrom: options.From, DateTo: options.To,
Status: model.SybSyncRunning, StartedAt: now.UTC().Format(model.TimeLayout),
}
if err := repository.CreateSybSyncRun(db, run); err != nil {
return "", err
}
return runID, nil
}
// FinishSybSyncRun 把同步报告持久化到对应记录。
func FinishSybSyncRun(db *sql.DB, runID string, report SyncReport) error {
status := model.SybSyncSucceeded
errorMessage := ""
if report.Err != nil {
status = model.SybSyncFailed
errorMessage = report.Err.Error()
if len([]rune(errorMessage)) > 500 {
errorMessage = string([]rune(errorMessage)[:500]) + "…"
}
}
finishedAt := report.FinishedAt
if finishedAt.IsZero() {
finishedAt = time.Now().UTC()
}
return repository.FinishSybSyncRun(db, model.SybSyncRun{
RunID: runID, Status: status, StockCount: report.StockCount,
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),
})
}
// InterruptRunningSybSyncRuns 收敛上次进程退出前未完成的同步记录。
func InterruptRunningSybSyncRuns(db *sql.DB, now time.Time) (int, error) {
return repository.InterruptRunningSybSyncRuns(db, now.UTC().Format(model.TimeLayout))
}
// ListSybSyncHistory 返回同步记录弹窗所需的分页视图。
func ListSybSyncHistory(db *sql.DB, page int) (*SybSyncHistoryResult, error) {
total, err := repository.CountSybSyncRuns(db)
if err != nil {
return nil, err
}
if page < 1 {
page = 1
}
totalPages := max(1, (total+sybSyncHistoryPageSize-1)/sybSyncHistoryPageSize)
if page > totalPages {
page = totalPages
}
runs, err := repository.ListSybSyncRuns(db, sybSyncHistoryPageSize, (page-1)*sybSyncHistoryPageSize)
if err != nil {
return nil, err
}
result := &SybSyncHistoryResult{Page: page, TotalPages: totalPages, Total: total}
for _, run := range runs {
statusText, statusClass := sybSyncRunStatusText(run.Status)
view := SybSyncRunView{
RunID: run.RunID, Username: run.Username,
DateRange: run.DateFrom + " ~ " + run.DateTo,
StatusText: statusText, StatusClass: statusClass,
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),
}
if run.CursorAdvanced {
view.CursorText = "已推进"
} else {
view.CursorText = "未推进"
}
result.Rows = append(result.Rows, view)
}
return result, nil
}
func sybSyncRunStatusText(status model.SybSyncRunStatus) (string, string) {
switch status {
case model.SybSyncRunning:
return "同步中", "status-running"
case model.SybSyncSucceeded:
return "成功", "status-success"
case model.SybSyncFailed:
return "失败", "status-failed"
case model.SybSyncInterrupted:
return "已中断", "status-interrupted"
default:
return "未知", ""
}
}
// Summary 组装状态条文案,格式见工单 #46「报告要说清楚」:
//
// 同步完成:日期范围 2026-08-09 ~ 2026-08-09,货运单 12 张,商品明细 27 条
// (新增 20,更新 7,跳过 0)
func (r SyncReport) Summary() string {
if r.Err != nil {
return "同步失败:" + r.Err.Error()
}
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 {
reasons = append(reasons, fmt.Sprintf("%s:%s", n.SybID, n.Reason))
}
msg += ";跳过/失败详情:" + strings.Join(reasons, ";")
}
return msg
}
// ---------- 同步互斥:一次只允许一个同步在跑 ----------
//
// `[必须]` 同步是长任务,不能阻塞 HTTP 请求线程直到结束;也不需要为此
// 引入后台协程池——本项目是单机内部工具,用一个互斥标志挡住重复点击
// 即可,见工单 #46。
var (
syncStateMu sync.Mutex
syncRunning bool
lastReport *SyncReport
)
// TryStartSybSync 尝试把"同步中"标志置上;已经在跑时返回 false。
func TryStartSybSync() bool {
syncStateMu.Lock()
defer syncStateMu.Unlock()
if syncRunning {
return false
}
syncRunning = true
return true
}
// FinishSybSync 同步结束(不管成功失败)后调用,记录最后一份报告并
// 放开互斥标志。
func FinishSybSync(r SyncReport) {
syncStateMu.Lock()
defer syncStateMu.Unlock()
syncRunning = false
rc := r
lastReport = &rc
}
// SybSyncStatus 是页面要显示的当前同步状态。
type SybSyncStatus struct {
Running bool
Report *SyncReport // 最近一次已经完成的同步报告,从没同步过是 nil
}
// GetSybSyncStatus 供页面渲染状态条用。
func GetSybSyncStatus() SybSyncStatus {
syncStateMu.Lock()
defer syncStateMu.Unlock()
return SybSyncStatus{Running: syncRunning, Report: lastReport}
}
// ---------- 同步主流程 ----------
// receiverFields 是顺运宝原始数据里属于个人信息的字段,落库前要剔掉。
//
// `[建议]` 见工单 #46 和 docs/admin/08-顺运宝接口.md §5:收件人姓名/
// 电话/地址做采购决策用不到,不入库。
var receiverFields = []string{"receiver", "receiverTel", "receiverAddr"}
func stripReceiverFields(m map[string]any) map[string]any {
if m == nil {
return nil
}
out := make(map[string]any, len(m))
for k, v := range m {
skip := false
for _, f := range receiverFields {
if k == f {
skip = true
break
}
}
if !skip {
out[k] = v
}
}
return out
}
// roundYuanToCent 把台币元换算成分,**先四舍五入再转整数**。
//
// `[必须]` 08 §5.1:直接截断浮点数会算错钱(239.0*100 在浮点下可能是
// 23899.999...),必须先 math.Round 再转 int64。
func roundYuanToCent(yuan float64) int64 {
return int64(math.Round(yuan * 100))
}
// sybDailyListResult 是一次单日列表分页尝试取得的数据。每次重试都重新创建,
// 防止把前一次已发生 offset 位移的 ID 混进下一次。
type sybDailyListResult struct {
stockByID map[int64]syb.StockRow
orderedIDs []int64
observedTotal int
}
// sybListDriftError 表示请求本身成功,但分页期间的列表没有形成稳定快照。
// 只有今天允许有限重试;历史日期遇到它仍然立即失败。
type sybListDriftError struct {
message string
}
func (e *sybListDriftError) Error() string { return e.message }
// loadSybDailyList 拉取一天的全部列表,并核对页长、分页前后总数及唯一 ID。
// 发生快照漂移时仍返回本次已经取得的合法 ID,供今天最后一次尝试在完整拉取
// 明细后安全 upsert;调用方不得因此把本次同步标记为成功或推进游标。
func loadSybDailyList(ctx context.Context, client *syb.Client, date string, pageSize, expectedTotal int) (sybDailyListResult, error) {
result := sybDailyListResult{
stockByID: make(map[int64]syb.StockRow, expectedTotal),
orderedIDs: make([]int64, 0, expectedTotal),
observedTotal: expectedTotal,
}
for start := 0; start < expectedTotal; start += pageSize {
pageIndex := start/pageSize + 1
rows, pageCount, err := client.ListPage(ctx, date, date, start, pageIndex, pageSize)
if err != nil {
return result, fmt.Errorf("拉取 %s 货运单列表第 %d 页失败(已获取 %d/%d 张,本次同步整体作废,"+
"下次会从同一个起始日期重新拉,靠 upsert 幂等不会重复计数): %w",
date, pageIndex, len(result.orderedIDs), expectedTotal, err)
}
for _, row := range rows {
if _, duplicate := result.stockByID[row.ID]; duplicate {
continue
}
result.stockByID[row.ID] = row
result.orderedIDs = append(result.orderedIDs, row.ID)
}
expectedPageCount := min(pageSize, expectedTotal-start)
if pageCount != expectedPageCount || len(rows) != expectedPageCount {
afterTotal, totalErr := client.ListTotal(ctx, date, date, pageSize)
if totalErr != nil {
return result, fmt.Errorf("分页异常后重新查询 %s 货运单总数失败: %w", date, totalErr)
}
result.observedTotal = afterTotal
return result, &sybListDriftError{message: fmt.Sprintf(
"%s 货运单列表第 %d 页不完整:预期 %d 行,实际 %d 行",
date, pageIndex, expectedPageCount, len(rows))}
}
}
afterTotal, err := client.ListTotal(ctx, date, date, pageSize)
if err != nil {
return result, fmt.Errorf("分页后重新查询 %s 货运单总数失败: %w", date, err)
}
result.observedTotal = afterTotal
if afterTotal != expectedTotal {
return result, &sybListDriftError{message: fmt.Sprintf(
"%s 货运单总数在分页期间从 %d 变为 %d", date, expectedTotal, afterTotal)}
}
if len(result.orderedIDs) != expectedTotal {
return result, &sybListDriftError{message: fmt.Sprintf(
"%s 货运单列表不完整:预期 %d 张,分页后只有 %d 个唯一 ID",
date, expectedTotal, len(result.orderedIDs))}
}
return result, nil
}
// RunSybSync 执行一次完整的顺运宝货运单同步。
//
// `[必须]` 调用方负责:
// 1. 只在 TryStartSybSync() 返回 true 时调用一次;
// 2. 调用结束后(不管成功失败)调用 FinishSybSync(report);
// 3. client 已经恢复了有效的登录会话(Cookie)。
//
// 本函数本身不检查会话是否有效——会话有效性判断属于"点同步"这一步的
// 前置检查(service.EnsureSybLoginNeeded),不属于同步本身;同步过程中
// 如果会话恰好失效,会从 syb.Client 的调用里冒出 syb.ErrSessionInvalid,
// 和其它错误一样按"中途失败"处理:不更新 last_synced_at,把原因写进报告。
func RunSybSync(ctx context.Context, db *sql.DB, client *syb.Client, cfg config.SybConfig, now time.Time) SyncReport {
return RunSybSyncWithOptions(ctx, db, client, cfg, now, SybSyncOptions{})
}
// RunSybSyncWithOptions 执行日期范围同步;空范围仅兼容旧版内部自动增量调用。
//
// 指定日期只有完整覆盖“本来应该自动同步的范围”时才推进 last_synced_at。
// 局部历史补拉只 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()}
allowedShops, err := repository.ListEnabledSybShopMappings(db)
if err != nil {
report.Err = fmt.Errorf("读取顺运宝允许店铺失败: %w", err)
report.FinishedAt = time.Now().UTC()
return report
}
if len(allowedShops) == 0 {
report.Err = fmt.Errorf("没有启用的顺运宝同步店铺,请先由管理员在“店铺管理”中配置并启用至少一个 SYB 店铺")
report.FinishedAt = time.Now().UTC()
return report
}
allowedNames := make([]string, 0, len(allowedShops))
for name := range allowedShops {
allowedNames = append(allowedNames, name)
}
sort.Strings(allowedNames)
report.ShopFilterHash = fmt.Sprintf("%x", sha256.Sum256([]byte(strings.Join(allowedNames, "\x00"))))
pageSize := cfg.PageSize
if pageSize <= 0 {
pageSize = 20
}
maxMatches := cfg.MaxMatches
if maxMatches <= 0 {
maxMatches = defaultSybMaxMatches
}
lastSyncedAt, hasCursor, err := repository.GetSybLastSyncedAt(db)
if err != nil {
report.Err = fmt.Errorf("读取上次同步进度失败: %w", err)
report.FinishedAt = time.Now().UTC()
return report
}
from, to := options.From, options.To
advanceCursor := !options.IsSpecified()
cursorAt := model.NowISO()
if options.IsSpecified() {
validated, validateErr := NewSybSyncOptions(options.From, options.To, now)
if validateErr != nil {
report.Err = validateErr
report.FinishedAt = time.Now().UTC()
return report
}
from, to = validated.From, validated.To
cursorAt, err = cursorISOForDate(to)
if err != nil {
report.Err = err
report.FinishedAt = time.Now().UTC()
return report
}
if !hasCursor {
// 统一日期入口第一次成功后,以操作员明确选择的结束日建立后续覆盖基线。
advanceCursor = true
} else {
lastTime, ok := model.ParseISO(lastSyncedAt)
if !ok {
report.Err = fmt.Errorf("上次同步时间 %q 解析失败", lastSyncedAt)
report.FinishedAt = time.Now().UTC()
return report
}
coveredDate := dateOf(lastTime)
// 只有范围衔接当前覆盖日期、且确实向后延伸时才推进。
// 历史补拉或跳过缺口的范围只 upsert 数据,不改变覆盖基线。
advanceCursor = from <= coveredDate && to > coveredDate
}
} else {
from, to, err = syncDateRange(lastSyncedAt, cfg.SyncFrom, now)
if err != nil {
report.Err = err
report.FinishedAt = time.Now().UTC()
return report
}
}
report.From, report.To = from, to
dates, err := splitDateRange(from, to)
if err != nil {
report.Err = err
report.FinishedAt = time.Now().UTC()
return report
}
type dailyPlan struct {
date string
total int
}
plans := make([]dailyPlan, 0, len(dates))
allTotal := 0
containsToday := false
today := dateOf(now)
for _, date := range dates {
if date == today {
containsToday = true
}
total, totalErr := client.ListTotal(ctx, date, date, pageSize)
if totalErr != nil {
report.Err = fmt.Errorf("查询 %s 货运单总数失败: %w", date, totalErr)
report.FinishedAt = time.Now().UTC()
return report
}
if total < 0 {
report.Err = fmt.Errorf("查询 %s 货运单总数返回负数 %d", date, total)
report.FinishedAt = time.Now().UTC()
return report
}
if total > maxMatches-allTotal {
report.Err = fmt.Errorf(
"日期范围 %s ~ %s 内货运单总数超过单次同步上限 %d,"+
"请缩小日期范围或在确认机器和网络容量后调大 max_matches", from, to, maxMatches)
report.FinishedAt = time.Now().UTC()
return report
}
allTotal += total
plans = append(plans, dailyPlan{date: date, total: total})
}
// 历史范围全为 0 时可以直接结束。范围包含今天时仍再做一次列表快照
// 校验,避免预检刚返回 0 就有新单进入而被当成完整成功。
if allTotal == 0 && !containsToday {
report.FinishedAt = time.Now().UTC()
if advanceCursor {
if err := repository.SetSybLastSyncedAt(db, cursorAt); err != nil {
report.Err = fmt.Errorf("更新同步进度失败: %w", err)
} else {
report.CursorAdvanced = true
}
}
return report
}
// ① 每天单独翻页,② 按 100 个一批取明细,③ 逐张货运单写库。
//
// `[必须]` 不放在一个大事务里——几千条明细的事务会长时间持锁;
// 按货运单为单位提交,失败了已成功的部分保留(下次重拉会 upsert
// 覆盖,幂等),见工单 #46。
//
// `[必须]` 中途失败(拉明细失败,或某张货运单写库失败)**立即停止**、
// 不更新 last_synced_at——已经成功写入的部分不回滚(它们本身是幂等
// 的),但"这次同步整体算成功"这件事不能发生,否则漏掉的单永远补不回来。
const detailBatch = 100
for _, plan := range plans {
if plan.total == 0 && plan.date != today {
continue
}
attempts := 1
if plan.date == today {
attempts = sybTodayListMaxAttempts
}
var listResult sybDailyListResult
var unstableTodayErr error
for attempt := 1; attempt <= attempts; attempt++ {
listResult, err = loadSybDailyList(ctx, client, plan.date, pageSize, plan.total)
if err == nil {
unstableTodayErr = nil
break
}
var driftErr *sybListDriftError
if plan.date != today || !errors.As(err, &driftErr) {
report.Err = fmt.Errorf("%w;本次同步停止且不推进游标", err)
report.FinishedAt = time.Now().UTC()
return report
}
unstableTodayErr = driftErr
if attempt == attempts {
break
}
if listResult.observedTotal < 0 {
report.Err = fmt.Errorf("重新查询 %s 货运单总数返回负数 %d",
plan.date, listResult.observedTotal)
report.FinishedAt = time.Now().UTC()
return report
}
otherDatesTotal := allTotal - plan.total
if listResult.observedTotal > maxMatches-otherDatesTotal {
report.Err = fmt.Errorf(
"今天货运单变化后日期范围 %s ~ %s 的总数超过单次同步上限 %d,"+
"已停止自动重试,请缩小日期范围或调大 max_matches",
from, to, maxMatches)
report.FinishedAt = time.Now().UTC()
return report
}
allTotal = otherDatesTotal + listResult.observedTotal
plan.total = listResult.observedTotal
}
otherDatesTotal := allTotal - plan.total
largestTodayCount := max(listResult.observedTotal, len(listResult.orderedIDs))
if largestTodayCount < 0 || largestTodayCount > maxMatches-otherDatesTotal {
report.Err = fmt.Errorf(
"今天货运单变化后日期范围 %s ~ %s 的总数超过单次同步上限 %d,"+
"未保存本次超限数据,请缩小日期范围或调大 max_matches",
from, to, maxMatches)
report.FinishedAt = time.Now().UTC()
return report
}
stockByID := listResult.stockByID
rawIDs := listResult.orderedIDs
report.StockCount += len(rawIDs)
orderedIDs := make([]int64, 0, len(rawIDs))
for _, id := range rawIDs {
if _, ok := sybShopID(allowedShops, stockByID[id].Raw, nil); !ok {
report.ShopSkipped++
continue
}
orderedIDs = append(orderedIDs, id)
}
report.AcceptedCount += len(orderedIDs)
for i := 0; i < len(orderedIDs); i += detailBatch {
end := i + detailBatch
if end > len(orderedIDs) {
end = len(orderedIDs)
}
batch := orderedIDs[i:end]
details, detailErr := client.DetailListByStock(ctx, batch)
if detailErr != nil {
report.Err = fmt.Errorf("拉取 %s 货运单明细失败(本次同步整体作废,"+
"已写入的数据保留,下次重拉会 upsert 覆盖): %w", plan.date, detailErr)
report.FinishedAt = time.Now().UTC()
return report
}
if completeErr := validateDetailBatch(batch, details); completeErr != nil {
report.Err = fmt.Errorf("%s 货运单明细不完整:%w;本次同步停止且不推进游标", plan.date, completeErr)
report.FinishedAt = time.Now().UTC()
return report
}
for _, d := range details {
stockRow := stockByID[d.ID]
shopID, ok := sybShopID(allowedShops, stockRow.Raw, d.Raw)
if !ok {
report.AcceptedCount--
report.ShopSkipped++
continue
}
if err := writeStockDetail(db, cfg.BaseURL, shopID, stockRow, d, &report); err != nil {
report.Err = fmt.Errorf("写入货运单 %s(id=%d)失败(本次同步整体作废,"+
"已写入的数据保留): %w", d.Code, d.ID, err)
report.FinishedAt = time.Now().UTC()
return report
}
}
}
if unstableTodayErr != nil {
report.Err = fmt.Errorf(
"%s 当天货运单在连续 %d 次分页期间仍有变化;最后一次取得原始货运单 %d 张,已保存其中允许店铺 %d 张的完整明细,"+
"本次未形成稳定快照且不推进游标,下次同步将继续覆盖当天:%w",
plan.date, sybTodayListMaxAttempts, len(rawIDs), len(orderedIDs), unstableTodayErr)
report.FinishedAt = time.Now().UTC()
return report
}
}
// ④ 全部成功,才更新 last_synced_at。
if advanceCursor {
if err := repository.SetSybLastSyncedAt(db, cursorAt); err != nil {
report.Err = fmt.Errorf("同步数据已全部写入,但更新同步进度失败,"+
"下次同步会重新拉这个日期范围(不会漏,但会重复拉一次): %w", err)
} else {
report.CursorAdvanced = true
}
}
report.FinishedAt = time.Now().UTC()
return report
}
// sybShopAllowed 用明细字段覆盖列表字段后再核对,防止列表通过但明细在同步
// 期间已变成其他店铺。detail 为空时只检查列表快照。
func sybShopID(allowed map[string]string, listRaw, detailRaw map[string]any) (string, bool) {
name := trimmedStringField(mergeRaw(listRaw, detailRaw), "shopName")
shopID, ok := allowed[name]
return shopID, ok
}
// validateDetailBatch 确认批量明细响应与请求 ID 一一对应。任何缺失、重复、
// 意外 ID 或空商品明细都会让同步失败,避免在数据不完整时推进游标。
func validateDetailBatch(requested []int64, details []syb.StockDetail) error {
wanted := make(map[int64]struct{}, len(requested))
for _, id := range requested {
wanted[id] = struct{}{}
}
seen := make(map[int64]struct{}, len(details))
for _, detail := range details {
if _, ok := wanted[detail.ID]; !ok {
return fmt.Errorf("响应包含未请求的货运单 id=%d", detail.ID)
}
if _, duplicate := seen[detail.ID]; duplicate {
return fmt.Errorf("响应重复返回货运单 id=%d", detail.ID)
}
if len(detail.Details) == 0 {
return fmt.Errorf("货运单 id=%d 没有返回商品明细", detail.ID)
}
seen[detail.ID] = struct{}{}
}
for _, id := range requested {
if _, ok := seen[id]; !ok {
return fmt.Errorf("响应缺少货运单 id=%d", id)
}
}
return nil
}
// writeStockDetail 把一张货运单的全部商品明细写进 syb_orders,
// 一张货运单一个事务(工单 #46「按货运单为单位提交」)。
func writeStockDetail(db *sql.DB, baseURL, shopID string, stockRow syb.StockRow, detail syb.StockDetail, report *SyncReport) error {
if len(detail.Details) == 0 {
return nil
}
tx, err := db.Begin()
if err != nil {
return fmt.Errorf("开始事务失败: %w", err)
}
defer tx.Rollback() // 已提交的事务再 Rollback 是空操作,安全
stockRaw := stripReceiverFields(mergeRaw(stockRow.Raw, detail.Raw))
for _, item := range detail.Details {
sybID := strconv.FormatInt(item.ID, 10)
if item.ProductQty <= 0 {
report.SkippedZero++
report.Notes = append(report.Notes, SkipNote{
SybID: sybID, Reason: fmt.Sprintf("数量为 %d,跳过(表结构要求 quantity > 0)", item.ProductQty),
})
continue
}
sybData, err := buildSybDataJSON(stockRaw, item.Raw)
if err != nil {
return fmt.Errorf("组装 syb_data 失败: %w", err)
}
order := model.SybOrder{
SybID: sybID,
OrderNo: detail.Code,
ShopID: shopID,
ShopName: trimmedStringField(stockRaw, "shopName"),
Title: item.ProductTitle,
ProductSpec: item.ProductSpec,
ShopeeGoodsID: strconv.FormatInt(item.ProductID, 10),
Quantity: item.ProductQty,
PriceTwdCent: roundYuanToCent(item.ProductPrice),
ImageURL: imageURLFromThumb(baseURL, item.ProductThumb),
SybData: sybData,
}
created, err := repository.UpsertSybOrder(tx, order)
if err != nil {
return err
}
if _, err := repository.UpsertSybShopeeSKU(tx, order.ShopeeGoodsID, order.ProductSpec, model.NowISO()); err != nil {
return fmt.Errorf("写入顺运宝规格观测失败: %w", err)
}
if created {
report.Created++
} else {
report.Updated++
}
report.DetailCount++
}
return tx.Commit()
}
// trimmedStringField 读取顺运宝原始响应中的文本字段。
// 类型不对、缺字段或只有空白时都按空值处理,不猜测或格式化上游数据。
func trimmedStringField(raw map[string]any, name string) string {
value, _ := raw[name].(string)
return strings.TrimSpace(value)
}
// mergeRaw 合并"货运单列表"和"货运单明细"两次响应里同一张货运单的
// 外层字段(不含 details),后者字段优先覆盖前者——明细响应更贴近
// "拉这批数据当下"的状态,见工单 #46 字段映射表 syb_data 的说明。
func mergeRaw(list, detail map[string]any) map[string]any {
out := make(map[string]any, len(list)+len(detail))
for k, v := range list {
out[k] = v
}
for k, v := range detail {
out[k] = v
}
return out
}
// buildSybDataJSON 组装落库的 syb_data:{"stock": ..., "detail": ...}。
// 嵌套两个 key 而不是拍平合并,是因为 stock 和 detail 两边都有名叫
// "id" 的字段,指的是完全不同的东西(货运单 id vs 明细行 id),
// 拍平会互相覆盖、审计时看不出原始结构。
func buildSybDataJSON(stockRaw, detailRaw map[string]any) (string, error) {
payload := map[string]any{"stock": stockRaw, "detail": detailRaw}
b, err := json.Marshal(payload)
if err != nil {
return "", err
}
return string(b), nil
}
// imageURLFromThumb 拼图片地址:{base_url}/api/p/file?id={productThumb},
// 见工单 #46 字段映射表。productThumb 是数字 ID,不是 URL;为 0 时
// 说明没有缩略图,返回空字符串(不拼一个指向 id=0 的坏链接)。
func imageURLFromThumb(baseURL string, productThumb int64) string {
if productThumb == 0 || baseURL == "" {
return ""
}
return strings.TrimRight(baseURL, "/") + "/api/p/file?id=" + strconv.FormatInt(productThumb, 10)
}
// ---------- 待登录客户端:验证码和登录共用同一个 Cookie Jar ----------
//
// `[必须]` 08 §3.2:验证码和登录必须用同一个 syb.Client(同一个 Cookie
// Jar),换客户端拿到的验证码就对不上。浏览器"取验证码图片"和"提交
// 登录表单"是两次独立的 HTTP 请求,Admin 侧要在这两次请求之间把同一个
// 客户端存住——本项目单机单操作员使用,用一个包级变量即可,
// 不需要按会话/用户区分。
var (
pendingLoginMu sync.Mutex
pendingLoginClient *syb.Client
)
// NewPendingSybLogin 为一次新的"取验证码 → 登录"流程创建客户端,
// 并存成"待登录"客户端,丢弃上一个(操作员点"换一张"验证码时,
// 上一张验证码本来就废了,不需要保留旧客户端)。
func NewPendingSybLogin(baseURL string) (*syb.Client, error) {
c, err := syb.New(baseURL)
if err != nil {
return nil, err
}
pendingLoginMu.Lock()
pendingLoginClient = c
pendingLoginMu.Unlock()
return c, nil
}
// PendingSybLoginClient 取出当前"待登录"客户端,没有则返回 nil——
// 调用方应该提示操作员先获取验证码。
func PendingSybLoginClient() *syb.Client {
pendingLoginMu.Lock()
defer pendingLoginMu.Unlock()
return pendingLoginClient
}
// ClearPendingSybLogin 清掉"待登录"客户端(登录成功或放弃时调用)。
func ClearPendingSybLogin() {
pendingLoginMu.Lock()
pendingLoginClient = nil
pendingLoginMu.Unlock()
}
// ---------- 会话有效性 ----------
// ErrSybLoginRequired 表示当前没有可用的顺运宝登录会话,
// 页面应该弹登录框,而不是报错。
var ErrSybLoginRequired = errors.New("顺运宝会话不存在或已过期,请重新登录")
// EnsureSybSession 检查本地缓存的顺运宝会话是否足够新鲜,够就把 Cookie
// 恢复进传入的 client;不够就返回 ErrSybLoginRequired(`errors.Is` 判断),
// 提示调用方走登录流程。
//
// `[决定]` 这里只做**本地**过期时间判断,不额外发一次
// GET /am/user/get 去问服务端"你还活着吗"——08 §3.5 描述的"网络故障不能
// 判定未登录"这条规则,在同步真正发起后、遇到任何一次 syb.ErrSessionInvalid
// 时同样会触发(syb.Client.do() 对所有 /am/** 接口都做了同一套分类),
// 不需要在这里再打一次专门的探测请求——省掉一次没有必要的网络往返,
// 也避免"探测请求本身超时"这种情况被误判成"未登录"。
func EnsureSybSession(db *sql.DB, client *syb.Client, username string, now time.Time) error {
cached, err := repository.GetSybSession(db, username)
if err != nil {
return fmt.Errorf("读取顺运宝会话缓存失败: %w", err)
}
if cached == nil {
return ErrSybLoginRequired
}
expiresAt, ok := model.ParseISO(cached.ExpiresAt)
if !ok || !now.Before(expiresAt) {
return ErrSybLoginRequired
}
if err := client.ImportCookiesJSON(cached.Cookies); err != nil {
return fmt.Errorf("恢复顺运宝会话失败: %w", err)
}
return nil
}
// ---------- 列表页 ----------
// SybOrderView 是列表页一行要显示的全部内容,已经格式化成字符串,
// 模板里不做判断和格式化,和其余四个模块的做法一致。
type SybOrderView struct {
SybID string
OrderNo string
ShopName string
Title string
ProductSpec string
ShopeeGoodsID string
Quantity int
PriceText string // "NT$239.00",和人民币价格一眼分得清
ImageURL string
Stage string
StageText string
StageHelp string
PddAssociationText string
ActionText string
NeedsAttention bool
CanCollect bool
CanPurchase bool
CanAIMatch bool
SelectionActions string
SelectionHint string
MappingSource string
MappingSourceText string
MappingSourceClass string
MappingSourceHelp string
DefaultUnitPrice string
DefaultTotalPrice string
MappedPddChoice string
MappingOptionKey string
ContextVersion string
UpdatedAt string
}
const (
SybStageSpecMissing = "spec_missing"
SybStagePddMissing = "pdd_missing"
SybStagePddPending = "pdd_pending"
SybStagePddCollecting = "pdd_collecting"
SybStagePddCollectingStale = "pdd_collecting_stale"
SybStagePddFailed = "pdd_failed"
SybStageMappingPending = "mapping_pending"
SybStageAIMatched = "ai_matched"
SybStagePurchaseReady = "purchase_ready"
SybStageTaskCreated = "task_created"
SybStagePurchaseCompleted = "purchase_completed"
SybStagePurchaseReview = "purchase_review"
SybStagePurchaseBlocked = "purchase_blocked"
)
// SybStageOption 是顺运宝处理阶段筛选项。
type SybStageOption struct {
Value string
Text string
}
func SybStageOptions() []SybStageOption {
return []SybStageOption{
{Value: "", Text: "全部"},
{Value: SybStageSpecMissing, Text: "顺运宝未提供规格"},
{Value: SybStagePddMissing, Text: "未关联 PDD"},
{Value: SybStagePddPending, Text: "PDD 待采集"},
{Value: SybStagePddCollecting, Text: "PDD 采集中"},
{Value: SybStagePddCollectingStale, Text: "PDD 采集中(超时)"},
{Value: SybStagePddFailed, Text: "PDD 采集失败"},
{Value: SybStageMappingPending, Text: "规格待匹配"},
{Value: SybStageAIMatched, Text: "AI规格匹配"},
{Value: SybStagePurchaseReady, Text: "可创建采购任务"},
{Value: SybStageTaskCreated, Text: "已创建采购任务"},
{Value: SybStagePurchaseCompleted, Text: "采购完成"},
{Value: SybStagePurchaseReview, Text: "采购待人工核对"},
{Value: SybStagePurchaseBlocked, Text: "采购数据异常"},
}
}
func ParseSybStage(raw string) string {
for _, option := range SybStageOptions() {
if raw == option.Value {
return raw
}
}
return ""
}
func SybStageLabel(stage string) string {
for _, option := range SybStageOptions() {
if option.Value == stage {
return option.Text
}
}
return ""
}
func sybStageFor(c repository.SybOrderContext) (stage, text, help, action string) {
switch {
case c.HasSucceededPurchaseTask:
return SybStagePurchaseCompleted, "采购完成", "采购任务已成功完成,请到采集采购页面查看 PDD 订单编号和下单时间。", "查看采购结果"
case c.HasManualReviewPurchaseTask:
return SybStagePurchaseReview, "采购待人工核对", "任务可能已经下单,必须先核对订单,禁止重新创建采购任务。", "核对采购订单"
case c.HasActiveTask:
return SybStageTaskCreated, "已创建采购任务", "已有未结束的采购任务,请到采集采购页面查看。", "查看采购任务"
case c.Order.SpecKey == "":
return SybStageSpecMissing, "采购数据异常:顺运宝未提供规格", "缺少规格原文,不能匹配或创建采购任务。", "核对缺失规格"
case c.PddGoodsID == "" || c.PddCollectStatus == "":
return SybStagePddMissing, "未关联 PDD", "下一步关联采购商品。", "关联 PDD"
case c.PddCollectStatus == string(model.CollectPending):
return SybStagePddPending, "PDD 待采集", "已关联 PDD 商品,尚未创建采集任务。", "创建采集任务"
case c.PddCollectStatus == string(model.CollectCollecting) && sybCollectingIsStale(c):
return SybStagePddCollectingStale, "PDD 采集中(超时)", "采集状态已超时或找不到有效任务,可以重新创建采集任务。", "重新创建采集任务"
case c.PddCollectStatus == string(model.CollectCollecting):
return SybStagePddCollecting, "PDD 采集中", "客户端正在采集规格和价格,请稍后刷新。", "查看采集状态"
case c.PddCollectStatus == string(model.CollectFailed):
help := "PDD 采集失败,可以重新创建采集任务。"
if c.PddCollectMsg != "" {
help += " 原因:" + c.PddCollectMsg
}
return SybStagePddFailed, "PDD 采集失败", help, "重新采集"
case !mappingIsValid(c):
return SybStageMappingPending, "规格待匹配", "PDD 数据已采集,下一步匹配采购规格。", "匹配规格"
case c.Order.Quantity <= 0:
return SybStagePurchaseBlocked, "采购数据异常", "顺运宝采购数量不是正数,不能创建采购任务。", "核对采购数量"
default:
return SybStagePurchaseReady, "可创建采购任务", "规格映射有效,可以创建采购任务。", "创建采购任务"
}
}
func sybCollectingIsStale(c repository.SybOrderContext) bool {
if !c.HasActiveCollectTask {
return true
}
updatedAt, ok := model.ParseISO(c.PddUpdatedAt)
return !ok || time.Since(updatedAt) > model.CollectStaleAfter
}
// SybListResult 是列表页要的全部数据。
type SybListResult struct {
Rows []SybOrderView
Total int
HasAny bool
IsFiltered bool
Stage string
Page int
TotalPages int
OrderNoQuery string
OrderNoCount int
MatchedOrderCount int
UnmatchedOrderNos []string
}
const (
maxSybOrderNoQueryLen = 4000
maxSybOrderNoCount = 100
maxSybOrderNoLen = 191
)
// ParseSybOrderNos 把采购员粘贴的订单号列变成去重后的精确查询值。
// 支持 Excel 换行/Tab,以及中英文逗号和分号;订单号本身不允许包含空白。
func ParseSybOrderNos(raw string) ([]string, error) {
raw = strings.TrimSpace(raw)
if raw == "" {
return nil, nil
}
if utf8.RuneCountInString(raw) > maxSybOrderNoQueryLen {
return nil, invalidFieldInput("order_no", "订单号输入不能超过 %d 个字符", maxSybOrderNoQueryLen)
}
parts := strings.FieldsFunc(raw, func(r rune) bool {
return unicode.IsSpace(r) || r == ',' || r == ',' || r == ';' || r == ';'
})
if len(parts) == 0 {
return nil, invalidFieldInput("order_no", "请输入至少一个有效订单号")
}
seen := make(map[string]struct{}, len(parts))
orderNos := make([]string, 0, len(parts))
for _, part := range parts {
orderNo := strings.TrimSpace(part)
if orderNo == "" {
continue
}
if utf8.RuneCountInString(orderNo) > maxSybOrderNoLen {
return nil, invalidFieldInput("order_no", "订单号 %q 超过 %d 个字符", orderNo, maxSybOrderNoLen)
}
if _, exists := seen[orderNo]; exists {
continue
}
seen[orderNo] = struct{}{}
orderNos = append(orderNos, orderNo)
if len(orderNos) > maxSybOrderNoCount {
return nil, invalidFieldInput("order_no", "一次最多搜索 %d 个不同订单号", maxSybOrderNoCount)
}
}
if len(orderNos) == 0 {
return nil, invalidFieldInput("order_no", "请输入至少一个有效订单号")
}
return orderNos, nil
}
// ListSybOrdersView 按筛选条件分页查货运单明细列表,翻成界面文字。
//
// `[必须]` 处理阶段由顺运宝规格键、PDD 关联、采集状态和规格映射实时推导,
// 不额外保存一份容易过期的状态字段。
func ListSybOrdersView(db *sql.DB, keyword, shop, stage string, page int) (*SybListResult, error) {
return ListSybOrdersViewWithPageSize(db, keyword, shop, stage, page, DefaultPageSize)
}
// ListSybOrdersViewWithPageSize 按白名单页容量查询顺运宝明细。
func ListSybOrdersViewWithPageSize(db *sql.DB, keyword, shop, stage string, page, pageSize int) (*SybListResult, error) {
pageSize = NormalizePageSize(pageSize)
shop = strings.TrimSpace(shop)
stage = ParseSybStage(stage)
orderNos, err := ParseSybOrderNos(keyword)
if err != nil {
hasAny, countErr := repository.CountSybOrdersTotal(db)
if countErr != nil {
return nil, countErr
}
return &SybListResult{
Rows: make([]SybOrderView, 0), HasAny: hasAny > 0, IsFiltered: true,
Stage: stage, Page: 1, TotalPages: 1, OrderNoQuery: keyword,
}, err
}
if stage == SybStagePddCollecting || stage == SybStagePddCollectingStale ||
stage == SybStageMappingPending || stage == SybStagePurchaseReady ||
stage == SybStageAIMatched ||
stage == SybStageTaskCreated || stage == SybStagePurchaseCompleted ||
stage == SybStagePurchaseReview || stage == SybStagePurchaseBlocked {
return listAdvancedSybStage(db, orderNos, shop, stage, page, pageSize)
}
filter := repository.SybOrderFilter{OrderNos: orderNos, Shop: shop, Stage: stage}
total, err := repository.CountSybOrders(db, filter)
if err != nil {
return nil, err
}
totalPages := TotalPages(total, pageSize)
page = ClampPage(page, totalPages)
offset := (page - 1) * pageSize
rows, err := repository.ListSybOrderContexts(db, filter, pageSize, offset)
if err != nil {
return nil, err
}
hasAny, err := repository.CountSybOrdersTotal(db)
if err != nil {
return nil, err
}
result := &SybListResult{
Rows: make([]SybOrderView, 0, len(rows)),
Total: total,
HasAny: hasAny > 0,
IsFiltered: len(orderNos) > 0 || shop != "" || stage != "",
Stage: stage,
Page: page,
TotalPages: totalPages,
OrderNoQuery: strings.Join(orderNos, ","),
OrderNoCount: len(orderNos),
}
for _, context := range rows {
result.Rows = append(result.Rows, sybOrderViewFor(context))
}
matchedOrderNos, err := repository.ListMatchedSybOrderNos(db, filter)
if err != nil {
return nil, err
}
setSybOrderMatchStats(result, orderNos, matchedOrderNos)
return result, nil
}
// listAdvancedSybStage 在 Service 解析最新 skus_json 后筛选实时推导阶段。
// 这些阶段不能只靠 SQL 判断:同一个 option key 是否仍存在,需要走唯一的
// OptionKey 规范化逻辑。当前同步量是百到千级,先保证采购判断正确;普通列表
// 和其余阶段仍在数据库分页。
func listAdvancedSybStage(db *sql.DB, orderNos []string, shop, stage string, page, pageSize int) (*SybListResult, error) {
contexts, err := repository.ListSybOrderContexts(db,
repository.SybOrderFilter{OrderNos: orderNos, Shop: shop}, -1, 0)
if err != nil {
return nil, err
}
filtered := make([]repository.SybOrderContext, 0)
for _, context := range contexts {
value, _, _, _ := sybStageFor(context)
matches := value == stage
if stage == SybStageAIMatched {
matches = value == SybStagePurchaseReady && context.MappingSource == "ai" && mappingIsValid(context)
}
if matches {
filtered = append(filtered, context)
}
}
total := len(filtered)
totalPages := TotalPages(total, pageSize)
page = ClampPage(page, totalPages)
start := (page - 1) * pageSize
if start > total {
start = total
}
end := start + pageSize
if end > total {
end = total
}
result := &SybListResult{Total: total, Stage: stage, Page: page, TotalPages: totalPages,
IsFiltered: true, Rows: make([]SybOrderView, 0, end-start),
OrderNoQuery: strings.Join(orderNos, ","), OrderNoCount: len(orderNos)}
hasAny, err := repository.CountSybOrdersTotal(db)
if err != nil {
return nil, err
}
result.HasAny = hasAny > 0
for _, context := range filtered[start:end] {
result.Rows = append(result.Rows, sybOrderViewFor(context))
}
matchedOrderNos := make([]string, 0, len(orderNos))
matchedSet := make(map[string]struct{}, len(orderNos))
for _, context := range filtered {
if _, exists := matchedSet[context.Order.OrderNo]; exists {
continue
}
matchedSet[context.Order.OrderNo] = struct{}{}
matchedOrderNos = append(matchedOrderNos, context.Order.OrderNo)
}
setSybOrderMatchStats(result, orderNos, matchedOrderNos)
return result, nil
}
func setSybOrderMatchStats(result *SybListResult, orderNos, matchedOrderNos []string) {
if len(orderNos) == 0 {
return
}
matched := make(map[string]struct{}, len(matchedOrderNos))
for _, orderNo := range matchedOrderNos {
matched[orderNo] = struct{}{}
}
result.MatchedOrderCount = len(matched)
result.UnmatchedOrderNos = make([]string, 0, len(orderNos)-len(matched))
for _, orderNo := range orderNos {
if _, exists := matched[orderNo]; !exists {
result.UnmatchedOrderNos = append(result.UnmatchedOrderNos, orderNo)
}
}
}
func sybOrderViewFor(context repository.SybOrderContext) SybOrderView {
o := context.Order
v := SybOrderView{SybID: o.SybID, OrderNo: o.OrderNo, ShopName: o.ShopName, Title: o.Title,
ProductSpec: o.ProductSpec, ShopeeGoodsID: o.ShopeeGoodsID,
Quantity: o.Quantity, ImageURL: o.ImageURL,
UpdatedAt: formatLocalTime(o.UpdatedAt)}
v.Stage, v.StageText, v.StageHelp, v.ActionText = sybStageFor(context)
v.PddAssociationText = sybPddAssociationText(context)
v.CanCollect = v.Stage == SybStagePddPending || v.Stage == SybStagePddFailed ||
v.Stage == SybStagePddCollectingStale
v.CanPurchase = v.Stage == SybStagePurchaseReady
v.CanAIMatch = v.Stage == SybStageMappingPending || v.Stage == SybStagePurchaseReady
actions := make([]string, 0, 3)
if v.CanCollect {
actions = append(actions, "collect")
}
if v.CanPurchase {
actions = append(actions, "purchase")
}
if v.CanAIMatch {
actions = append(actions, "ai")
}
v.SelectionActions = strings.Join(actions, " ")
switch {
case len(actions) > 0:
labels := make([]string, 0, len(actions))
if v.CanCollect {
labels = append(labels, "创建采集")
}
if v.CanPurchase {
labels = append(labels, "创建采购")
}
if v.CanAIMatch {
labels = append(labels, "AI规格匹配")
}
v.SelectionHint = "可用于:" + strings.Join(labels, "、")
default:
v.SelectionHint = "当前阶段不能批量操作:" + v.StageHelp
}
if context.MappingOptionKey != "" {
v.MappingSource = context.MappingSource
if v.MappingSource == "" {
v.MappingSource = "manual"
}
v.MappingSourceText = mappingSourceText(v.MappingSource)
v.MappingSourceClass = "source-" + v.MappingSource
v.MappingSourceHelp = context.MappingReason
if !mappingIsValid(context) {
v.MappingSourceText += "(已失效)"
}
}
if v.CanPurchase {
v.MappingOptionKey = context.MappingOptionKey
v.ContextVersion = mappingContextVersion(context)
choice, err := findPddChoice(context.PddSkusJSON, context.MappingOptionKey)
if err == nil && choice != nil {
v.MappedPddChoice = choice.Label
if choice.HasPrice {
v.DefaultUnitPrice = fmt.Sprintf("%d.%02d", choice.PriceCent/100, choice.PriceCent%100)
if total, totalErr := CalculateOrderPriceLimitCent(choice.PriceCent, o.Quantity); totalErr == nil {
v.DefaultTotalPrice = fmt.Sprintf("%d.%02d", total/100, total%100)
}
}
}
}
v.NeedsAttention = v.Stage != SybStagePddCollecting && v.Stage != SybStageTaskCreated &&
v.Stage != SybStagePurchaseCompleted
if o.PriceTwdCent > 0 {
v.PriceText = fmt.Sprintf("NT$%.2f", float64(o.PriceTwdCent)/100)
} else {
v.PriceText = placeholder
}
if v.Title == "" {
v.Title = placeholder
}
if v.ProductSpec == "" {
v.ProductSpec = placeholder
}
return v
}
// GetSybOrderView 按稳定明细 ID 读取一行完整视图,供任务来源深链接在分页之外
// 准备采购确认数据。查不到返回 (nil, nil)。
func GetSybOrderView(db *sql.DB, sybID string) (*SybOrderView, error) {
context, err := repository.GetSybOrderContext(db, strings.TrimSpace(sybID))
if err != nil || context == nil {
return nil, err
}
view := sybOrderViewFor(*context)
return &view, nil
}
// SybProcessingDetail 是顺运宝“下一步”弹窗第一阶段需要的上下文。
type SybProcessingDetail struct {
SybID string
OrderNo string
Title string
ProductSpec string
ShopeeGoodsID string
Quantity int
Stage string
StageText string
StageHelp string
ShopeeExists bool
PddGoodsID string
PddURL string
PddAssociationText string
CollectStatus string
CollectMsg string
CanCollect bool
PddDimensionNames []string
PddChoices []PddOptionChoice
ContextVersion string
RecommendationNotice string
MappingValid bool
HasActiveTask bool
MappingSource string
MappingSourceText string
MappingSourceClass string
MappingSourceHelp string
}
func GetSybProcessingDetail(db *sql.DB, sybID string) (*SybProcessingDetail, error) {
context, err := repository.GetSybOrderContext(db, strings.TrimSpace(sybID))
if err != nil || context == nil {
return nil, err
}
o := context.Order
d := &SybProcessingDetail{
SybID: o.SybID, OrderNo: o.OrderNo, Title: o.Title, ProductSpec: o.ProductSpec,
ShopeeGoodsID: o.ShopeeGoodsID,
Quantity: o.Quantity, ShopeeExists: context.ShopeeExists,
PddGoodsID: context.PddGoodsID, PddURL: context.PddGoodsURL,
PddAssociationText: sybPddAssociationText(*context),
CollectStatus: context.PddCollectStatus, CollectMsg: context.PddCollectMsg,
}
d.Stage, d.StageText, d.StageHelp, _ = sybStageFor(*context)
d.CanCollect = context.PddGoodsID != "" &&
(d.Stage == SybStagePddPending || d.Stage == SybStagePddFailed ||
d.Stage == SybStagePddCollectingStale)
d.HasActiveTask = context.HasActiveTask
if context.MappingOptionKey != "" {
d.MappingSource = context.MappingSource
if d.MappingSource == "" {
d.MappingSource = "manual"
}
d.MappingSourceText = mappingSourceText(d.MappingSource)
d.MappingSourceClass = "source-" + d.MappingSource
d.MappingSourceHelp = context.MappingReason
if !mappingIsValid(*context) {
d.MappingSourceText += "(已失效)"
}
}
if context.PddCollectStatus == string(model.CollectCollected) && context.PddSkusJSON != "" {
choices, keys, names, parseErr := pddOptionChoices(context.PddSkusJSON)
if parseErr == nil {
d.PddDimensionNames = names
match := rankSpecChoices(o.ProductSpec, choices, keys, names)
d.PddChoices = match.Choices
d.RecommendationNotice = match.Notice
for i := range d.PddChoices {
d.PddChoices[i].Selected = d.PddChoices[i].Key == context.MappingOptionKey ||
(context.MappingOptionKey == "" && d.PddChoices[i].Key == match.PreselectOptionKey)
if d.PddChoices[i].Selected {
d.MappingValid = true
}
}
}
}
d.ContextVersion = mappingContextVersion(*context)
return d, nil
}
// sybPddAssociationText 只在当前蝦皮商品关联能找到有效、未删除的 PDD 商品时显示。
// PDD 关联的唯一事实来源是 shopee_products;顺运宝明细不复制这份关系。
func sybPddAssociationText(context repository.SybOrderContext) string {
if context.PddGoodsID == "" || context.PddCollectStatus == "" {
return ""
}
return "使用蝦皮商品的 PDD 关联:" + context.PddGoodsID + ",无需重复输入链接"
}
// SaveSybLoginSession 登录成功后把会话缓存进数据库。
//
// `[必须]` 缓存写失败不能让已经登录的会话失效——本函数把错误原样
// 返回,由调用方决定"写失败了但登录已经成功,要不要继续走后面的同步",
// 不在这里吞掉错误也不在这里替调用方做决定。
func SaveSybLoginSession(db *sql.DB, client *syb.Client, username string, expiresAt time.Time) error {
cookiesJSON, err := client.ExportCookiesJSON()
if err != nil {
return fmt.Errorf("导出顺运宝会话 Cookie 失败: %w", err)
}
return repository.SaveSybSession(db, username, cookiesJSON, expiresAt.UTC().Format(model.TimeLayout))
}