Files
cmautobuy/admin/service/syb.go
T

1224 lines
44 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"
"strconv"
"strings"
"sync"
"time"
"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
)
// 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 通常预填最近三天;如果覆盖游标落后,则优先从游标
// 当天连续补齐。缺口超过 31 天时只选最早一段,成功后下一次继续。
func DefaultSybSyncRange(db *sql.DB, now time.Time) (SybSyncDefaults, error) {
today, err := time.Parse(dateLayout, dateOf(now))
if err != nil {
return SybSyncDefaults{}, fmt.Errorf("计算顺运宝今天日期失败: %w", err)
}
recentStart := today.AddDate(0, 0, -2)
defaults := SybSyncDefaults{From: recentStart.Format(dateLayout), To: today.Format(dateLayout)}
lastSyncedAt, found, err := repository.GetSybLastSyncedAt(db)
if err != nil {
return SybSyncDefaults{}, err
}
if !found {
return defaults, nil
}
lastTime, ok := model.ParseISO(lastSyncedAt)
if !ok {
return SybSyncDefaults{}, fmt.Errorf("上次同步时间 %q 解析失败", lastSyncedAt)
}
coveredDate, err := time.Parse(dateLayout, dateOf(lastTime))
if err != nil || !coveredDate.Before(recentStart) {
return defaults, nil
}
segmentEnd := coveredDate.AddDate(0, 0, maxSpecifiedSyncDays-1)
defaults.From = coveredDate.Format(dateLayout)
if segmentEnd.Before(today) {
defaults.To = segmentEnd.Format(dateLayout)
defaults.Warning = fmt.Sprintf(
"存在较长的未同步区间,已先选择最早 %d 天(%s 至 %s);本段成功后请继续同步下一段。",
maxSpecifiedSyncDays, defaults.From, defaults.To)
return defaults, nil
}
defaults.To = today.Format(dateLayout)
defaults.Warning = fmt.Sprintf("检测到上次同步停在 %s,已自动扩展开始日期以补齐缺口。", defaults.From)
return defaults, 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 // 拉到的货运单数
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,
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)",
run.StockCount, 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)",
r.From, r.To, r.StockCount, 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))
}
// 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()}
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
for _, date := range dates {
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})
}
if allTotal == 0 {
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 {
continue
}
stockByID := map[int64]syb.StockRow{}
orderedIDs := make([]int64, 0, plan.total)
for start := 0; start < plan.total; start += pageSize {
pageIndex := start/pageSize + 1
rows, pageCount, listErr := client.ListPage(ctx, plan.date, plan.date, start, pageIndex, pageSize)
if listErr != nil {
report.Err = fmt.Errorf("拉取 %s 货运单列表第 %d 页失败(已获取 %d/%d 张,本次同步整体作废,"+
"下次会从同一个起始日期重新拉,靠 upsert 幂等不会重复计数): %w",
plan.date, pageIndex, len(orderedIDs), plan.total, listErr)
report.FinishedAt = time.Now().UTC()
return report
}
expectedPageCount := min(pageSize, plan.total-start)
if pageCount != expectedPageCount || len(rows) != expectedPageCount {
report.Err = fmt.Errorf("%s 货运单列表第 %d 页不完整:预期 %d 行,实际 %d 行,"+
"本次同步停止且不推进游标", plan.date, pageIndex, expectedPageCount, len(rows))
report.FinishedAt = time.Now().UTC()
return report
}
for _, row := range rows {
if _, dup := stockByID[row.ID]; dup {
continue
}
stockByID[row.ID] = row
orderedIDs = append(orderedIDs, row.ID)
}
}
afterTotal, totalErr := client.ListTotal(ctx, plan.date, plan.date, pageSize)
if totalErr != nil {
report.Err = fmt.Errorf("分页后重新查询 %s 货运单总数失败: %w", plan.date, totalErr)
report.FinishedAt = time.Now().UTC()
return report
}
if afterTotal != plan.total {
report.Err = fmt.Errorf("%s 货运单总数在分页期间从 %d 变为 %d,"+
"为防止 offset 分页漏单,本次同步停止且不推进游标", plan.date, plan.total, afterTotal)
report.FinishedAt = time.Now().UTC()
return report
}
if len(orderedIDs) != plan.total {
report.Err = fmt.Errorf("%s 货运单列表不完整:预期 %d 张,分页后只有 %d 个唯一 ID,"+
"本次同步停止且不推进游标", plan.date, plan.total, len(orderedIDs))
report.FinishedAt = time.Now().UTC()
return report
}
report.StockCount += 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]
if err := writeStockDetail(db, cfg.BaseURL, 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
}
}
}
}
// ④ 全部成功,才更新 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
}
// 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 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,
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 created {
report.Created++
} else {
report.Updated++
}
report.DetailCount++
}
return tx.Commit()
}
// 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
Title string
ProductSpec string
ShopeeGoodsID string
Quantity int
PriceText string // "NT$239.00",和人民币价格一眼分得清
ImageURL string
Stage string
StageText string
StageHelp string
ActionText string
NeedsAttention bool
CanPurchase bool
DefaultMaxPrice string
MappedPddChoice string
MappingOptionKey string
ContextVersion string
UpdatedAt string
}
const (
SybStageSpecMissing = "spec_missing"
SybStagePddMissing = "pdd_missing"
SybStagePddPending = "pdd_pending"
SybStagePddCollecting = "pdd_collecting"
SybStagePddFailed = "pdd_failed"
SybStageMappingPending = "mapping_pending"
SybStagePurchaseReady = "purchase_ready"
SybStageTaskCreated = "task_created"
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: SybStagePddFailed, Text: "PDD 采集失败"},
{Value: SybStageMappingPending, Text: "规格待匹配"},
{Value: SybStagePurchaseReady, Text: "可创建采购任务"},
{Value: SybStageTaskCreated, 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.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):
return SybStagePddCollecting, "PDD 采集中", "客户端正在采集规格和价格,请稍后刷新。", "查看采集状态"
case c.PddCollectStatus == string(model.CollectFailed):
help := "PDD 采集失败,可以重新创建采集任务。"
if c.PddCollectMsg != "" {
help += " 原因:" + c.PddCollectMsg
}
return SybStagePddFailed, "PDD 采集失败", help, "重新采集"
case c.HasActiveTask:
return SybStageTaskCreated, "已创建采购任务", "已有未结束的采购任务,请到采集采购页面查看。", "查看采购任务"
case !mappingIsValid(c):
return SybStageMappingPending, "规格待匹配", "PDD 数据已采集,下一步匹配采购规格。", "匹配规格"
case c.Order.Quantity <= 0:
return SybStagePurchaseBlocked, "采购数据异常", "顺运宝采购数量不是正数,不能创建采购任务。", "核对数据"
default:
return SybStagePurchaseReady, "可创建采购任务", "规格映射有效,可以创建采购任务。", "创建采购任务"
}
}
// SybListResult 是列表页要的全部数据。
type SybListResult struct {
Rows []SybOrderView
Total int
HasAny bool
IsFiltered bool
Stage string
Page int
TotalPages int
}
// ListSybOrdersView 按筛选条件分页查货运单明细列表,翻成界面文字。
//
// `[必须]` 处理阶段由顺运宝规格键、PDD 关联、采集状态和规格映射实时推导,
// 不额外保存一份容易过期的状态字段。
func ListSybOrdersView(db *sql.DB, keyword, stage string, page int) (*SybListResult, error) {
stage = ParseSybStage(stage)
if stage == SybStageMappingPending || stage == SybStagePurchaseReady ||
stage == SybStageTaskCreated || stage == SybStagePurchaseBlocked {
return listAdvancedSybStage(db, keyword, stage, page)
}
filter := repository.SybOrderFilter{Keyword: keyword, Stage: stage}
total, err := repository.CountSybOrders(db, filter)
if err != nil {
return nil, err
}
totalPages := TotalPages(total)
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: strings.TrimSpace(keyword) != "" || stage != "",
Stage: stage,
Page: page,
TotalPages: totalPages,
}
for _, context := range rows {
result.Rows = append(result.Rows, sybOrderViewFor(context))
}
return result, nil
}
// listAdvancedSybStage 在 Service 解析最新 skus_json 后筛选映射相关阶段。
// 这四个阶段不能只靠 SQL 判断:同一个 option key 是否仍存在,需要走唯一的
// OptionKey 规范化逻辑。当前同步量是百到千级,先保证采购判断正确;普通列表
// 和其余阶段仍在数据库分页。
func listAdvancedSybStage(db *sql.DB, keyword, stage string, page int) (*SybListResult, error) {
contexts, err := repository.ListSybOrderContexts(db,
repository.SybOrderFilter{Keyword: keyword}, -1, 0)
if err != nil {
return nil, err
}
filtered := make([]repository.SybOrderContext, 0)
for _, context := range contexts {
value, _, _, _ := sybStageFor(context)
if value == stage {
filtered = append(filtered, context)
}
}
total := len(filtered)
totalPages := TotalPages(total)
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)}
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))
}
return result, nil
}
func sybOrderViewFor(context repository.SybOrderContext) SybOrderView {
o := context.Order
v := SybOrderView{SybID: o.SybID, OrderNo: o.OrderNo, 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.CanPurchase = v.Stage == SybStagePurchaseReady
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.DefaultMaxPrice = fmt.Sprintf("%d.%02d", choice.PriceCent/100, choice.PriceCent%100)
}
}
}
v.NeedsAttention = v.Stage != SybStagePddCollecting && v.Stage != SybStageTaskCreated
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
}
// 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
CollectStatus string
CollectMsg string
CanCollect bool
PddDimensionNames []string
PddChoices []PddOptionChoice
ContextVersion string
RecommendationNotice string
MappingValid bool
HasActiveTask bool
}
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,
CollectStatus: context.PddCollectStatus, CollectMsg: context.PddCollectMsg,
}
d.CanCollect = context.PddGoodsID != "" &&
(context.PddCollectStatus == string(model.CollectPending) ||
context.PddCollectStatus == string(model.CollectFailed))
d.HasActiveTask = context.HasActiveTask
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)
d.Stage, d.StageText, d.StageHelp, _ = sybStageFor(*context)
return d, nil
}
// 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))
}