// 顺运宝货运单同步的编排逻辑:算日期范围、翻页拉列表和明细、 // 字段映射、落库统计、汇总报告。 // // 改动前必读 admin/AGENTS.md:本层不认识 *gin.Context,也不拼 SQL—— // 那些分别在 handler/web/others.go 和 repository/syb.go。 // // 接口契约见 docs/admin/08-顺运宝接口.md,工单见 #46。三条最容易出事的规则: // 1. shopee_sku_id 绝不能被同步写入/覆盖——这条已经在 // repository.UpsertSybOrder 的 SQL 层面保证,本文件不需要、 // 也不允许再传一份"新的" shopee_sku_id 进去。 // 2. 增量必须从"上次同步日期当天"重新拉,不是第二天,见 syncDateRange。 // 3. 中途失败不更新 last_synced_at,见 RunSybSync 的最后一步。 package service import ( "context" "database/sql" "encoding/json" "errors" "fmt" "math" "strconv" "strings" "sync" "time" "cmautobuy/admin/config" "cmautobuy/admin/model" "cmautobuy/admin/repository" "cmautobuy/admin/syb" ) 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), // `[必须]` 不设置 ShopeeSKUID——顺运宝没有这个值, // repository.UpsertSybOrder 也不会用它覆盖已有的匹配结果。 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 ShopeeSKUID string Quantity int PriceText string // "NT$239.00",和人民币价格一眼分得清 ImageURL string MatchText string // "已匹配" / "待匹配" Matched bool UpdatedAt string } // SybListResult 是列表页要的全部数据。 type SybListResult struct { Rows []SybOrderView Total int HasAny bool IsFiltered bool Page int TotalPages int } // ListSybOrdersView 按筛选条件分页查货运单明细列表,翻成界面文字。 // // `[必须]` 匹配状态是**算出来的**(shopee_sku_id 是否非空), // 不是存的字段,见工单 #46——本工单不做规格匹配,这里只负责如实 // 显示"有没有"。 func ListSybOrdersView(db *sql.DB, keyword string, page int) (*SybListResult, error) { filter := repository.SybOrderFilter{Keyword: keyword} 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.ListSybOrders(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) != "", Page: page, TotalPages: totalPages, } for _, o := range rows { v := SybOrderView{ SybID: o.SybID, OrderNo: o.OrderNo, Title: o.Title, ProductSpec: o.ProductSpec, ShopeeGoodsID: o.ShopeeGoodsID, ShopeeSKUID: o.ShopeeSKUID, Quantity: o.Quantity, ImageURL: o.ImageURL, UpdatedAt: formatLocalTime(o.UpdatedAt), } if o.PriceTwdCent > 0 { v.PriceText = fmt.Sprintf("NT$%.2f", float64(o.PriceTwdCent)/100) } else { v.PriceText = placeholder } if o.Title == "" { v.Title = placeholder } if o.ProductSpec == "" { v.ProductSpec = placeholder } if o.ShopeeSKUID != "" { v.Matched = true v.MatchText = "已匹配" } else { v.MatchText = "待匹配" } result.Rows = append(result.Rows, v) } return result, 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)) }