feat: 未选择店铺时同步全部店铺商品 (#26)

This commit is contained in:
QiuSW
2026-09-29 17:18:16 +08:00
parent 1a2b6feaf6
commit 309e3eeb11
10 changed files with 563 additions and 42 deletions
+28 -17
View File
@@ -50,8 +50,9 @@ type App struct {
log *logx.Logger
probe downloader.ProbeFunc
videoTaskMu sync.Mutex
videoTask *task.Runner
videoTaskMu sync.Mutex
videoTask *task.Runner
productSyncMu sync.Mutex
}
// NewApp 创建应用对象。真正的初始化在 startup 里做。
@@ -818,37 +819,47 @@ func (a *App) ExportLogs() (string, error) {
// ---------------------------------------------------------------- 商品操作
// DownloadProductData 通过 erpgo 拉取商品列表到本地(需求 R1)。
func (a *App) DownloadProductData(platformShopID string) error {
func (a *App) DownloadProductData(platformShopID string) (erpgo.SyncResult, error) {
if a.db == nil {
return erpgo.BridgeError(fmt.Errorf("数据库未就绪"))
return erpgo.SyncResult{}, erpgo.BridgeError(fmt.Errorf("数据库未就绪"))
}
if platformShopID == "" {
return &erpgo.Error{Code: "INVALID_ARGUMENT", Message: "请先选择店铺"}
if !a.productSyncMu.TryLock() {
return erpgo.SyncResult{}, &erpgo.Error{Code: "SYNC_IN_PROGRESS", Message: "已有商品同步正在进行,请等待完成"}
}
defer a.productSyncMu.Unlock()
client, err := erpgo.NewClient(a.cfg.ERPGo, nil)
if err != nil {
return err
return erpgo.SyncResult{}, err
}
ctx := a.ctx
if ctx == nil {
ctx = context.Background()
}
a.log.Info("开始下载所选店铺的商品数据")
products, diagnoses, err := client.DownloadAllProducts(ctx, platformShopID, func(current, total int) {
a.log.Info("已拉取 %d/%d 页", current, total)
a.log.Info("开始同步商品数据")
result, err := client.SyncProductData(ctx, a.db, platformShopID, func(p erpgo.SyncProgress) {
if p.Current == 0 {
a.log.Info("正在同步第 %d/%d 个店铺", p.ShopIndex, p.TotalShops)
} else {
a.log.Info("第 %d/%d 个店铺已拉取 %d/%d 页", p.ShopIndex, p.TotalShops, p.Current, p.Pages)
}
})
if err != nil {
a.log.Error("下载商品数据失败:%s", erpgo.LogSummary(err))
return err
return result, err
}
updatedAt := time.Now().Format("2006-01-02 15:04:05")
if err := a.db.SyncProducts(products, diagnoses, updatedAt); err != nil {
a.log.Error("保存商品与诊断失败:%s", erpgo.LogSummary(err))
return erpgo.BridgeError(err)
for _, failed := range result.Failures {
a.log.Warn("第 %d 个店铺同步失败:%s", failed.ShopIndex, erpgo.LogSummary(failed.Error))
}
a.log.Success("商品数据下载完成,共拉取 %d 条", len(products))
return nil
if result.StopError != nil {
a.log.Warn("商品同步已停止:%s", erpgo.LogSummary(result.StopError))
}
if result.FailedShops > 0 || result.StopError != nil {
a.log.Warn("商品同步结束:成功 %d,失败 %d,未处理 %d 个店铺,共保存 %d 条商品", result.SucceededShops, result.FailedShops, result.SkippedShops, result.ProductCount)
} else {
a.log.Success("商品数据同步完成:%d 个店铺,共保存 %d 条商品", result.SucceededShops, result.ProductCount)
}
return result, nil
}
// StartVideoTask 启动选中商品的批量取视频任务。
+1 -1
View File
@@ -38,7 +38,7 @@ func TestERPGoLiveReadOnlySync(t *testing.T) {
if len(shops) == 0 {
t.Skip("no shops available for product integration")
}
if err := a.DownloadProductData(shops[0].PlatformShopID); err != nil {
if _, err := a.DownloadProductData(shops[0].PlatformShopID); err != nil {
t.Fatalf("live product pagination or persistence failed: %v", err)
}
page, err := a.db.ListProducts(store.ProductQuery{PlatformShopID: shops[0].PlatformShopID, Page: 1, PageSize: 1})
+37 -3
View File
@@ -8,11 +8,45 @@ import (
"path/filepath"
"strings"
"testing"
"time"
"cmsp/internal/config"
"cmsp/internal/erpgo"
"cmsp/internal/store"
)
func TestAppProductSyncRejectsOverlappingRequests(t *testing.T) {
started, release := make(chan struct{}), make(chan struct{})
a := newSyncTestApp(t, func(w http.ResponseWriter, r *http.Request) {
close(started)
<-release
queryResponse(w, 200, appProductPage(1, 1), "")
})
completed := make(chan error, 1)
go func() { _, err := a.DownloadProductData("demo-shop"); completed <- err }()
defer func() {
close(release)
select {
case err := <-completed:
if err != nil {
t.Errorf("first sync failed: %v", err)
}
case <-time.After(5 * time.Second):
t.Error("first sync did not finish")
}
}()
select {
case <-started:
case <-time.After(5 * time.Second):
t.Fatal("first sync did not start")
}
_, err := a.DownloadProductData("")
if e, ok := err.(*erpgo.Error); !ok || e.Code != "SYNC_IN_PROGRESS" {
t.Fatalf("overlapping sync was accepted: %v", err)
}
assertOldSyncData(t, a)
}
func newSyncTestApp(t *testing.T, handler http.HandlerFunc) *App {
t.Helper()
server := httptest.NewServer(handler)
@@ -80,7 +114,7 @@ func TestAppQueryAndAtomicProductSync(t *testing.T) {
}
queryResponse(w, 200, appProductPage(1, pages), "")
})
err := a.DownloadProductData("demo-shop")
_, err := a.DownloadProductData("demo-shop")
if failSecondPage {
if err == nil || !strings.Contains(err.Error(), "HHH_UPSTREAM_ERROR") {
t.Fatalf("missing stable error: %v", err)
@@ -105,7 +139,7 @@ func TestAppDiagnosisWriteFailureRollsBackSync(t *testing.T) {
if _, err := a.db.DB().Exec(`CREATE TRIGGER fail_sync BEFORE INSERT ON product_diagnoses BEGIN SELECT RAISE(ABORT,'fictional write failure'); END`); err != nil {
t.Fatal(err)
}
err := a.DownloadProductData("demo-shop")
_, err := a.DownloadProductData("demo-shop")
if err == nil || !strings.Contains(err.Error(), "LOCAL_SYNC_FAILED") {
t.Fatalf("missing write error: %v", err)
}
@@ -146,7 +180,7 @@ func TestAppMissingQuerySettingsRetainsDataAndUploadSettings(t *testing.T) {
a.cfg.ERPGo = config.ERPGoConfig{}
a.cfg.Huohanhan.Account = "fictional-upload-account"
a.cfg.Huohanhan.Password = "fictional-upload-password"
err := a.DownloadProductData("demo-shop")
_, err := a.DownloadProductData("demo-shop")
if err == nil || !strings.Contains(err.Error(), "ERPGo_NOT_CONFIGURED") {
t.Fatal("missing configuration not reported")
}
+13 -4
View File
@@ -2,12 +2,20 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Architecture-and-Code-Map
wiki_url: https://git.ilapage.cn/chengma/cmsp/wiki/Architecture-and-Code-Map.-
wiki_revision: c070f41cdae8b495e381dfd571500274697a0677
synchronized_at: 2026-09-29T08:52:08Z
wiki_revision: b6f38cfa6fd1abde9bf864e070843eceaa8f1651
synchronized_at: 2026-09-29T09:17:31Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
## 全部店铺商品同步(2026-09-29,#26)
DownloadProductData(platformShopID) 返回 internal/erpgo.SyncResult 和 error。指定店铺只查询该店铺;空/空白参数先实时 ListShops 并更新缓存,再按返回顺序串行执行 DownloadAllProducts,每店每页200。internal/erpgo/sync.go 编排查询和逐店事务,不依赖界面当前筛选条件。
每店完整分页校验、内部ID去重后调用 store.SyncProducts,商品和诊断同一事务;失败不提交该店部分页。已成功店铺保留,普通店铺错误继续;跨店同内部ID拒绝后一个店铺,避免覆盖已成功数据。API_KEY_INVALID、HHH_AUTH_FAILED、RATE_LIMITED、REQUEST_CANCELLED、NETWORK_ERROR、SERVICE_UNAVAILABLE、INTERNAL_ERROR 全局停止,不自动重试,剩余店铺记为未处理。
结果 allShops、totalShops、succeededShops、failedShops、skippedShops、productCount 来自实际处理/提交;failures 包含店铺序号、标识、名称和安全 error;stopError 为全局停止原因。初始化与单店失败走 error 通道,全店开始后的部分失败返回汇总。App.productSyncMu 在Go侧拒绝重叠同步,返回 SYNC_IN_PROGRESS;前端禁用重复操作,显示可关闭汇总和可展开失败详情。日志只写序号、计数和稳定错误摘要,不输出店铺标识、名称或原始响应。
## 本地列表分页条数(2026-09-29)
商品列表 n-pagination 启用 show-size-picker,提供每页20/50/100/200/500条;初始/重置仍20。changePageSize 更新 query.pageSize、清空旧勾选,调用 search 回到第一页,其他筛选条件保留。仅此次页大小切换明确清空勾选,普通搜索/翻页继续原选择行为。
@@ -45,7 +53,8 @@ store.ListProducts 接受1—500的 pageSize,空/非正数/大于500仍回退2
RefreshShops → internal/erpgo.Client.ListShops
→ GET /api/v1/integrations/huohanhan/shops
→ internal/store.ReplaceShops(同一事务替换本地店铺缓存)
DownloadProductData → internal/erpgo.Client.DownloadAllProducts
DownloadProductData → internal/erpgo.Client.SyncProductData(选店单店,空选实时全部店铺)
→ 每店 internal/erpgo.Client.DownloadAllProducts
→ GET /api/v1/integrations/huohanhan/products(指定店铺,串行分页)
→ 完整校验并按内部商品 id 去重
→ internal/store.SyncProducts(商品和诊断同一事务)
@@ -54,7 +63,7 @@ DownloadProductData → internal/erpgo.Client.DownloadAllProducts
主要入口是 app.go 的 RefreshShops、DownloadProductData;internal/erpgo/client.go 负责白名单转换、稳定错误码及分页完整性;internal/store/product.go 的 SyncProducts 使用与诊断仓储共用的事务辅助函数,不改变 SQLite 表结构。
查询只读取货憨憨现有数据,不触发 Shopee 同步;不自动回退直连。任何请求/校验失败不返回可落库的部分数据;数据库失败回滚商品与诊断。下载、上传、视频记录和断点状态保留,不依据列表删除商品。internal/huohanhan 仍用于现有视频上传,原客户端代码保留用于回归与回退。
查询只读取货憨憨现有数据,不触发 Shopee 同步;不自动回退直连。每店请求/校验失败不返回该店可落库的部分数据;数据库失败回滚该店商品与诊断。全店模式保留此前已成功店铺。下载、上传、视频记录和断点状态保留,不依据列表删除商品。internal/huohanhan 仍用于现有视频上传,原客户端代码保留用于回归与回退。
## 当前实现状态
+12 -4
View File
@@ -2,12 +2,20 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Business-Rules-and-Glossary
wiki_url: https://git.ilapage.cn/chengma/cmsp/wiki/Business-Rules-and-Glossary.-
wiki_revision: 8a8fc50cb1c9163b6819cc834e6511c27ed12bf8
synchronized_at: 2026-09-29T08:52:08Z
wiki_revision: 34592c1fc27459e30f215ef053c3d0c603526505
synchronized_at: 2026-09-29T09:17:36Z
<!-- gitea-wiki-mirror:end -->
# 业务规则与术语
## 未选择店铺的同步范围(2026-09-29,#26)
“下载数据”选店时同步所选店铺,不选店时按 erpgo 当前返回的全部店铺逐店串行同步。下载的是在售商品数据,不会自动下载/上传视频或触发 Shopee 同步;列表的缺视频、下载状态与分页筛选不缩小远端同步范围。
全店同步按店铺提交:每店完整分页成功后保存商品和诊断,保留原下载/上传/视频状态、last_error、视频记录及断点。失败店铺原数据保留,已成功店铺不回滚;合法空列表不删除旧商品。实时店铺列表成功后独立替换缓存,即使后续商品查询失败;店铺查询本身失败则保留原缓存,不开始商品查询。
普通店铺错误记录后继续,认证、Key失效、限流、取消、网络或服务不可用等全局错误停止。界面分别报告成功、失败、未处理店铺和实际保存商品数,不把部分成功当全部成功。同店分页按内部ID去重;跨店ID冲突不覆盖已成功店铺。Go侧拒绝重叠商品同步,排除原因后可选失败店铺手动重试。
## 本地商品列表每页条数(2026-09-29)
分页栏可选20、50、100、200、500条,默认及重置20。切换每页条数后回到第一页、保留店铺/诊断/下载/上传等筛选、清空旧勾选;列表总数仍为全部匹配记录数。500条是显示范围,不表示自动批量处理500条,批量操作仍需勾选。
@@ -42,7 +50,7 @@ MTOP 搜索方式、签名与请求参数保持现状。商品原链接以实际
- id 为货憨憨内部商品 ID;itemId 为 Shopee 商品号;platformShopId 为店铺号。按 id 去重,重复记录采用后取得的数据及诊断,不以 itemId 代替内部 ID。
- videoDiagnosis=missing 表示明确诊断“缺少视频”;ok 只表示未出现该诊断,包括无诊断,不证明真实视频存在或已在 Shopee 生效。前端使用稳定枚举,不按中文诊断文本分支。
- fetchedAt 是货憨憨读取时间,不是 Shopee 更新时间;createTime 保留上游字符串。上游分页不是一致性快照,去重不能保证读取期间没有遗漏。
- 请求、响应或分页校验失败,保留原商品、诊断、店铺缓存和任务状态;不将部分结果视为成功。完整分页成功后商品和诊断同一事务落库,任一写入失败全部回滚。
- 每店请求、响应或分页校验失败,保留该店原商品、诊断和任务状态;不将该店部分页视为成功。该店完整分页成功后商品和诊断同一事务落库,任一写入失败该店事务回滚。全店模式保留此前已成功店铺;店铺查询成功独立刷新缓存。
- 成功同步只更新商品基础字段与诊断;保留 video_status、download_status、upload_status、视频记录、last_error 与任务断点,不依据查询列表删除商品,不引入额外工作流状态重置。
- 查询不触发 Shopee 同步、上传、修改或删除。店铺成功刷新可按既有规则替换缓存;合法空商品集合不删除本地商品。
- 新 API Key、X-API-Key 与配置原文禁止日志输出;YAML 解析错误不透传可能包含字段原值的解析器文本。
@@ -52,7 +60,7 @@ MTOP 搜索方式、签名与请求参数保持现状。商品原链接以实际
| 术语 | 含义 |
|---|---|
| 货憨憨 / HHH | 内部使用的跨境电商 ERP,是商品数据的来源,也是写回 Shopee 的唯一通道。程序访问的是它的 HTTP 接口 |
| Shopee 商品 | 在货憨憨中登记、最终展示在 Shopee 店铺的在线商品。本项目只处理指定店铺的商品 |
| Shopee 商品 | 在货憨憨中登记、最终展示在 Shopee 店铺的在线商品。本项目支持所选店铺或全部店铺的在售商品数据同步 |
| 商品主图 | 货憨憨商品的第一张图片,用作淘宝以图搜的输入 |
| 拍立淘 / 以图搜 | 淘宝按图片搜索同款商品的能力。本项目通过 MTOP 接口调用,不使用网页交互 |
| MTOP | 淘宝的 H5 接口网关。请求需要 `appKey`、毫秒时间戳、由 `_m_h5_tk` 派生的 token 和 MD5 签名 |
@@ -2,12 +2,18 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Local-Development-and-Verification
wiki_url: https://git.ilapage.cn/chengma/cmsp/wiki/Local-Development-and-Verification.-
wiki_revision: 4219925d9609a8c4c6343744a08af090774fc96f
synchronized_at: 2026-09-29T08:52:09Z
wiki_revision: 1b97e699e41e235f146e2f26c703e5246cc9eba5
synchronized_at: 2026-09-29T09:17:38Z
<!-- gitea-wiki-mirror:end -->
# 本地开发与验证
## 全店铺同步验证(2026-09-29,#26)
internal/erpgo/sync_test.go 使用模拟HTTP与临时SQLite验证最新店铺列表、串行店铺/分页、每店完整提交、同店ID去重、分页及本地事务失败继续、全局错误/取消停止、空店铺/列表失败、跨店ID冲突,以及诊断、下载/上传状态和已上传视频记录保留。app_sync_test.go 验证单店请求及Go侧重叠同步拒绝。
验证命令:go test ./...、go test -race ./...、go vet ./...;Windows 使用 Wails 构建生成绑定及前端资源。模拟测试及构建不等于真实全店铺联调或 Windows WebView2/Narrator/高对比度/DPI验收;本次没有对真实全部店铺落库,没有淘宝访问或外部写操作。默认真实联调测试仍只处理首个店铺的临时数据库,不自动扩大为全部店铺。
## 每页500条边界验证(2026-09-29)
internal/store/product_test.go 使用501条虚构匹配记录和1条不匹配记录,验证20/50/100/200/500五档有效、无效值回退20、500条末页1条、随后空页、分页不重叠与总数一致。前端使用已安装Naive UI分页控件的showSizePicker/pageSizes/onUpdate:pageSize契约;changePageSize 保留筛选、清空勾选并回到第一页。默认/重置为20。
+13 -3
View File
@@ -2,12 +2,22 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Troubleshooting
wiki_url: https://git.ilapage.cn/chengma/cmsp/wiki/Troubleshooting
wiki_revision: c91a79baab19786e491b46e8aceb2308a2ea49d5
synchronized_at: 2026-09-29T07:39:14Z
wiki_revision: c94bf22f0832f7665e9c866f95ee7f85bdc10bdf
synchronized_at: 2026-09-29T09:18:00Z
<!-- gitea-wiki-mirror:end -->
# 故障排查
## 全店铺同步与部分成功(2026-09-29,#26)
不选店铺点击“下载数据”先获取最新全部店铺,逐店串行查询;耗时随店铺及商品量增加,运行日志显示店铺序号和页数。完成汇总来自实际保存数量,与当前缺视频/未下载筛选的列表总数不同。
“同步部分完成”可展开失败店铺看错误;已成功店铺已保存,失败及未处理店铺原商品和任务状态保留。分页或落库失败不写该店部分结果;店铺列表成功时缓存可以刷新。不要清库、删除视频或重置上传状态,排除原因后选失败店铺手动重试。
API_KEY_INVALID、HHH_AUTH_FAILED、RATE_LIMITED、REQUEST_CANCELLED、NETWORK_ERROR、SERVICE_UNAVAILABLE、INTERNAL_ERROR 停止剩余店铺,不逐店重试全局错误。SHOP_ACCESS_DENIED、HHH_UPSTREAM_ERROR/TIMEOUT、INVALID_RESPONSE 和 LOCAL_SYNC_FAILED 等店铺错误记录后继续。不同店铺同内部ID冲突返回 INVALID_RESPONSE,不覆盖先成功的数据。
SYNC_IN_PROGRESS 表示已有商品同步,请等待结束;Go侧拒绝重叠请求。空店铺列表显示“没有可同步的店铺”,不删除旧商品。店铺查询失败则不开始商品查询,原缓存保留。
## 图搜成功但商品链接无法使用(2026-09-29)
程序优先使用 MTOP 返回的 `auctionURL`,不会丢弃原查询参数后重新按 ID 拼接。单个非法链接候选会跳过,不影响其他合法候选;如果所有候选都不可用,报“淘宝以图搜返回了无效商品链接”,检查经过脱敏的字段结构是否发生变化;只允许淘宝/天猫商品域名、`/item.htm` 和与候选匹配的唯一 `id`。不能关闭校验、盲信广告跳转域名或把完整响应/参数写入日志。缺少链接字段的旧响应仍兼容按 ID 打开。
@@ -28,7 +38,7 @@ synchronized_at: 2026-09-29T07:39:14Z
## erpgo 店铺与商品同步排查(2026-09-28)
先查看稳定 errorCode、HTTP 状态及 requestId。界面按错误码显示提示,不按中文上游消息分支,也不输出原始响应。所有查询失败保留本地缓存和任务;不要清库或删视频来修复连接问题。
先查看稳定 errorCode、HTTP 状态及 requestId。界面按错误码显示提示,不按中文上游消息分支,也不输出原始响应。失败店铺保留原商品和任务;全店模式保留已成功店铺,实时店铺列表成功后缓存可刷新。不要清库或删视频来修复连接问题。
| errorCode | 处理 |
|---|---|
+49 -8
View File
@@ -497,18 +497,20 @@ async function run(action, fn) {
}
async function downloadProductData() {
if (!query.value.platformShopId) {
message.warning('请先在上面选择一个店铺')
return
}
const shopID = String(query.value.platformShopId || '').trim()
// 商品多的店铺要拉几十页,耗时十几秒到几分钟。
// 必须先给出可见反馈,否则按钮看起来像坏了。
running.value = '下载数据'
message.info('正在下载商品数据,进度见「运行日志」')
productSyncResult.value = null
message.info(shopID ? '正在下载所选店铺数据,进度见「运行日志」' : '正在逐店铺下载全部店铺数据,进度见「运行日志」')
try {
await DownloadProductData(query.value.platformShopId)
const result = await DownloadProductData(shopID)
productSyncResult.value = result
if (result.allShops) await loadShops()
await search(false)
message.success(`下载完成,本店铺共 ${total.value} 条商品`)
if (result.failedShops || result.stopError) message.warning(productSyncSummary.value)
else if (!result.totalShops) message.info(productSyncSummary.value)
else message.success(productSyncSummary.value)
} catch (err) {
message.error(`下载数据失败:${queryErrorMessage(err)}`)
} finally {
@@ -519,6 +521,14 @@ async function downloadProductData() {
// 正在跑的操作名。非空时按钮禁用并显示转圈,
// 否则长耗时操作期间界面毫无反馈,使用者会以为点了没用。
const running = ref('')
const productSyncResult = ref(null)
const productSyncSummary = computed(() => {
const result = productSyncResult.value
if (!result) return ''
if (!result.totalShops) return '没有可同步的店铺'
const state = result.stopError ? '同步已停止' : result.failedShops ? '同步部分完成' : '同步完成'
return `${state}:成功 ${result.succeededShops}/${result.totalShops} 个店铺,失败 ${result.failedShops},未处理 ${result.skippedShops},实际保存 ${result.productCount} 条商品`
})
const selectedCount = computed(() => checkedIds.value.length)
@@ -573,6 +583,7 @@ function queryErrorMessage(err) {
REQUEST_CANCELLED: '查询已取消',
INVALID_RESPONSE: 'erpgo 返回的数据或分页不符合契约,请联系维护者',
LOCAL_SYNC_FAILED: '本地同步失败,请查看运行日志',
SYNC_IN_PROGRESS: '已有商品同步正在进行,请等待完成',
}
const text = messages[detail.errorCode] || '同步失败,请联系维护者'
const reference = detail.requestId ? `(请求编号:${detail.requestId})` : ''
@@ -797,7 +808,7 @@ onUnmounted(() => {
<n-button
type="primary"
:loading="running === '下载数据'"
:disabled="taskRunning || (running !== '' && running !== '下载数据')"
:disabled="taskRunning || running !== ''"
@click="downloadProductData"
>
下载数据
@@ -822,6 +833,30 @@ onUnmounted(() => {
</div>
</div>
<n-alert
v-if="productSyncResult"
:type="productSyncResult.failedShops || productSyncResult.stopError ? 'warning' : productSyncResult.totalShops ? 'success' : 'info'"
:title="productSyncSummary"
closable
aria-live="polite"
@close="productSyncResult = null"
>
<p v-if="productSyncResult.stopError || productSyncResult.failedShops">
成功店铺的数据已保存,失败及未处理店铺的原商品数据和任务状态保留。
</p>
<p v-if="productSyncResult.stopError">
停止原因:{{ queryErrorMessage(JSON.stringify(productSyncResult.stopError)) }}
</p>
<details v-if="productSyncResult.failures.length">
<summary>查看失败店铺({{ productSyncResult.failures.length }})</summary>
<ul class="sync-failures">
<li v-for="failure in productSyncResult.failures" :key="failure.platformShopId">
{{ failure.shopName || failure.platformShopId }}:{{ queryErrorMessage(JSON.stringify(failure.error)) }}
</li>
</ul>
</details>
</n-alert>
<div v-if="taskRunning" class="task-bar" aria-live="polite">
<div class="task-main">
<div class="task-current" :title="taskProgress.currentName || taskProgress.current">
@@ -1044,6 +1079,12 @@ onUnmounted(() => {
padding: 0 24px;
flex-shrink: 0;
}
.sync-failures {
margin: 8px 0 0;
padding-left: 20px;
max-height: 160px;
overflow: auto;
}
.muted {
color: var(--muted);
font-size: 12px;
+115
View File
@@ -0,0 +1,115 @@
package erpgo
import (
"context"
"errors"
"strings"
"time"
"cmsp/internal/store"
)
type SyncFailure struct {
ShopIndex int `json:"shopIndex"`
PlatformShopID string `json:"platformShopId"`
ShopName string `json:"shopName"`
Error *Error `json:"error"`
}
// SyncResult 的计数来自实际提交,不使用界面当前筛选后的商品数量。
type SyncResult struct {
AllShops bool `json:"allShops"`
TotalShops int `json:"totalShops"`
SucceededShops int `json:"succeededShops"`
FailedShops int `json:"failedShops"`
SkippedShops int `json:"skippedShops"`
ProductCount int `json:"productCount"`
Failures []SyncFailure `json:"failures"`
StopError *Error `json:"stopError,omitempty"`
}
type SyncProgress struct {
ShopIndex, TotalShops, Current, Pages int
}
// SyncProductData 按店铺串行查询、逐店铺原子保存,失败不提交该店铺的部分页。
// 全店模式保留已成功店铺,并返回失败/未处理范围;不自动重试失败请求。
func (c *Client) SyncProductData(ctx context.Context, db *store.Store, shopID string, onProgress func(SyncProgress)) (SyncResult, error) {
shopID = strings.TrimSpace(shopID)
result := SyncResult{AllShops: shopID == "", Failures: make([]SyncFailure, 0)}
shops := []store.Shop{{PlatformShopID: shopID}}
if result.AllShops {
var err error
shops, err = c.ListShops(ctx)
if err != nil {
return result, err
}
if ctx.Err() != nil {
return result, failure("REQUEST_CANCELLED", 0, "")
}
if err := db.ReplaceShops(shops, time.Now().Format("2006-01-02 15:04:05")); err != nil {
return result, BridgeError(err)
}
}
result.TotalShops = len(shops)
committedIDs := make(map[string]bool)
for i, shop := range shops {
if ctx.Err() != nil {
result.StopError = failure("REQUEST_CANCELLED", 0, "")
result.SkippedShops = len(shops) - i
break
}
if onProgress != nil {
onProgress(SyncProgress{ShopIndex: i + 1, TotalShops: len(shops)})
}
products, diagnoses, err := c.DownloadAllProducts(ctx, shop.PlatformShopID, func(current, pages int) {
if onProgress != nil {
onProgress(SyncProgress{ShopIndex: i + 1, TotalShops: len(shops), Current: current, Pages: pages})
}
})
if ctx.Err() != nil {
err = failure("REQUEST_CANCELLED", 0, "")
}
if err == nil {
for _, product := range products {
if committedIDs[product.ID] {
err = failure("INVALID_RESPONSE", 200, "")
break
}
}
}
if err == nil {
err = db.SyncProducts(products, diagnoses, time.Now().Format("2006-01-02 15:04:05"))
}
if err != nil {
var safe *Error
errors.As(BridgeError(err), &safe)
result.FailedShops++
result.Failures = append(result.Failures, SyncFailure{ShopIndex: i + 1, PlatformShopID: shop.PlatformShopID, ShopName: shop.ShopName, Error: safe})
if !result.AllShops {
return result, safe
}
if stopProductSync(safe.Code) {
result.StopError = safe
result.SkippedShops = len(shops) - i - 1
break
}
continue
}
for _, product := range products {
committedIDs[product.ID] = true
}
result.SucceededShops++
result.ProductCount += len(products)
}
return result, nil
}
func stopProductSync(code string) bool {
switch code {
case "API_KEY_INVALID", "HHH_AUTH_FAILED", "RATE_LIMITED", "REQUEST_CANCELLED", "NETWORK_ERROR", "SERVICE_UNAVAILABLE", "INTERNAL_ERROR":
return true
default:
return false
}
}
+287
View File
@@ -0,0 +1,287 @@
package erpgo
import (
"context"
"encoding/json"
"net/http"
"path/filepath"
"reflect"
"strconv"
"strings"
"sync"
"testing"
"cmsp/internal/store"
)
func syncTestDB(t *testing.T) *store.Store {
t.Helper()
db, err := store.Open(filepath.Join(t.TempDir(), "fictional.db"))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = db.Close() })
if err := db.SyncProducts([]store.Product{{ID: "a-product", PlatformShopID: "a", ItemName: "old-title"}}, map[string][]store.Diagnosis{"a-product": {{Type: "old-diagnosis"}}}, "old-time"); err != nil {
t.Fatal(err)
}
if err := db.UpdateProductStatus("a-product", store.VideoFound, store.DownloadDone, store.UploadDone, "old-error"); err != nil {
t.Fatal(err)
}
if err := db.ReplaceVideos("a-product", []store.Video{{SourceItem: "fictional-source", Status: store.VideoStatusUploaded, RemoteURL: "https://example.invalid/video.mp4"}}, "old-time"); err != nil {
t.Fatal(err)
}
if err := db.ReplaceShops([]store.Shop{{ID: "cached", PlatformShopID: "cached", ShopName: "cached"}}, "old-time"); err != nil {
t.Fatal(err)
}
return db
}
func writeSyncShops(w http.ResponseWriter, ids ...string) {
items := make([]Shop, 0, len(ids))
for _, id := range ids {
items = append(items, Shop{ID: "fictional-" + id, PlatformShopID: id, ShopName: "fictional " + id, Platform: "0"})
}
writeResponse(w, 200, map[string]any{"source": "huohanhan", "fetchedAt": "2026-09-28T10:00:00+08:00", "items": items}, "")
}
func syncTestPage(shop string, current, pages int, id string) productPage {
p := sampleProduct(id)
p.PlatformShopID, p.ItemName = shop, "new-title"
page := samplePage(current, p)
page.PlatformShopID, page.Pages = shop, pages
if pages == 1 {
page.Total = 1
}
return page
}
func assertSyncProduct(t *testing.T, db *store.Store, unchanged bool) {
t.Helper()
p, found, err := db.GetProduct("a-product")
d, de := db.ListDiagnoses("a-product")
v, ve := db.ListVideos("a-product")
wantTitle, wantDiagnosis := "new-title", "缺少视频"
if unchanged {
wantTitle, wantDiagnosis = "old-title", "old-diagnosis"
}
if err != nil || !found || de != nil || ve != nil || p.ItemName != wantTitle || p.VideoStatus != store.VideoFound || p.DownloadStatus != store.DownloadDone || p.UploadStatus != store.UploadDone || p.LastError != "old-error" || len(d) != 1 || d[0].Type != wantDiagnosis || len(v) != 1 || v[0].Status != store.VideoStatusUploaded || v[0].RemoteURL != "https://example.invalid/video.mp4" {
t.Fatalf("product data/state/video changed unexpectedly: %+v", p)
}
if unchanged && p.SyncedAt != "old-time" {
t.Fatal("failed shop was committed")
}
}
func TestSyncAllShopsSerialPagesAndPartialFailure(t *testing.T) {
for _, fail := range []bool{false, true} {
t.Run(strconv.FormatBool(fail), func(t *testing.T) {
db := syncTestDB(t)
var mu sync.Mutex
var requests []string
client := testClient(t, func(w http.ResponseWriter, r *http.Request) {
if r.Method != "GET" || r.Header.Get("X-API-Key") != fictionalKey {
t.Error("unexpected request")
}
mu.Lock()
requests = append(requests, r.URL.Path+"?"+r.URL.RawQuery)
mu.Unlock()
if strings.HasSuffix(r.URL.Path, "/shops") {
writeSyncShops(w, "a", "b")
return
}
shop := r.URL.Query().Get("platformShopId")
current, _ := strconv.Atoi(r.URL.Query().Get("current"))
if shop == "a" {
if current == 2 && fail {
writeResponse(w, 502, nil, "HHH_UPSTREAM_ERROR")
return
}
if current == 2 {
assertSyncProduct(t, db, true)
}
writeResponse(w, 200, syncTestPage(shop, current, 2, "a-product"), "")
return
}
if shop != "b" {
t.Error("did not use latest shop list")
}
assertSyncProduct(t, db, fail)
writeResponse(w, 200, syncTestPage(shop, 1, 1, "b-product"), "")
})
result, err := client.SyncProductData(context.Background(), db, " ", nil)
if err != nil || !result.AllShops || result.TotalShops != 2 || result.SkippedShops != 0 || result.StopError != nil {
t.Fatalf("unexpected result: %+v, %v", result, err)
}
wantSucceeded, wantFailed := 2, 0
if fail {
wantSucceeded, wantFailed = 1, 1
}
if result.SucceededShops != wantSucceeded || result.ProductCount != wantSucceeded || result.FailedShops != wantFailed || len(result.Failures) != wantFailed {
t.Fatalf("incorrect actual counts: %+v", result)
}
if fail {
assertCode(t, result.Failures[0].Error, "HHH_UPSTREAM_ERROR")
}
assertSyncProduct(t, db, fail)
mu.Lock()
defer mu.Unlock()
prefix := "/api/v1/integrations/huohanhan/"
want := []string{prefix + "shops?", prefix + "products?current=1&platformShopId=a&size=200", prefix + "products?current=2&platformShopId=a&size=200", prefix + "products?current=1&platformShopId=b&size=200"}
if !reflect.DeepEqual(requests, want) {
t.Fatalf("requests not serial: %v", requests)
}
encoded, _ := json.Marshal(result)
if strings.Contains(string(encoded), fictionalKey) {
t.Fatal("raw upstream error leaked")
}
})
}
}
func TestSyncGlobalErrorsStopButShopErrorsContinue(t *testing.T) {
cases := []struct {
code string
status int
stop bool
}{
{"API_KEY_INVALID", 401, true}, {"HHH_AUTH_FAILED", 502, true}, {"RATE_LIMITED", 429, true}, {"SERVICE_UNAVAILABLE", 503, true}, {"INTERNAL_ERROR", 500, true}, {"NETWORK_ERROR", 502, true}, {"REQUEST_CANCELLED", 502, true},
{"SHOP_ACCESS_DENIED", 403, false}, {"INVALID_ARGUMENT", 400, false}, {"HHH_UPSTREAM_ERROR", 502, false}, {"HHH_UPSTREAM_TIMEOUT", 504, false}, {"INVALID_RESPONSE", 502, false},
}
for _, tc := range cases {
t.Run(tc.code, func(t *testing.T) {
db := syncTestDB(t)
client := testClient(t, func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/shops") {
writeSyncShops(w, "a", "b")
return
}
if r.URL.Query().Get("platformShopId") == "a" {
writeResponse(w, tc.status, nil, tc.code)
return
}
if tc.stop {
t.Error("global failure did not stop requests")
}
writeResponse(w, 200, syncTestPage("b", 1, 1, "b-product"), "")
})
result, err := client.SyncProductData(context.Background(), db, "", nil)
if err != nil || result.FailedShops != 1 || len(result.Failures) != 1 {
t.Fatalf("missing partial report: %+v %v", result, err)
}
assertCode(t, result.Failures[0].Error, tc.code)
if tc.stop {
if result.StopError == nil || result.SkippedShops != 1 || result.SucceededShops != 0 {
t.Fatalf("incorrect stop: %+v", result)
}
} else if result.StopError != nil || result.SkippedShops != 0 || result.SucceededShops != 1 {
t.Fatalf("shop error stopped batch: %+v", result)
}
assertSyncProduct(t, db, true)
})
}
}
func TestSyncEmptyShopsAndListFailure(t *testing.T) {
for _, fail := range []bool{false, true} {
t.Run(strconv.FormatBool(fail), func(t *testing.T) {
db := syncTestDB(t)
client := testClient(t, func(w http.ResponseWriter, r *http.Request) {
if !strings.HasSuffix(r.URL.Path, "/shops") {
t.Error("unexpected products request")
}
if fail {
writeResponse(w, 401, nil, "API_KEY_INVALID")
} else {
writeSyncShops(w)
}
})
result, err := client.SyncProductData(context.Background(), db, "", nil)
if fail {
assertCode(t, err, "API_KEY_INVALID")
} else if err != nil || result.TotalShops != 0 || result.ProductCount != 0 {
t.Fatalf("empty shops failed: %+v %v", result, err)
}
shops, err := db.ListShops()
if err != nil || (fail && (len(shops) != 1 || shops[0].PlatformShopID != "cached")) || (!fail && len(shops) != 0) {
t.Fatal("shop cache did not preserve/refresh correctly")
}
assertSyncProduct(t, db, true)
})
}
}
func TestSyncCancellationDoesNotCommitPartialShop(t *testing.T) {
db := syncTestDB(t)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
client := testClient(t, func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/shops") {
writeSyncShops(w, "a", "b")
return
}
if r.URL.Query().Get("platformShopId") != "a" || r.URL.Query().Get("current") != "1" {
t.Error("request continued after cancellation")
}
writeResponse(w, 200, syncTestPage("a", 1, 2, "a-product"), "")
})
result, err := client.SyncProductData(ctx, db, "", func(p SyncProgress) {
if p.Current == 1 {
cancel()
}
})
if err != nil || result.FailedShops != 1 || result.SkippedShops != 1 || result.ProductCount != 0 {
t.Fatalf("cancellation failed: %+v %v", result, err)
}
assertCode(t, result.StopError, "REQUEST_CANCELLED")
assertSyncProduct(t, db, true)
}
func TestSyncCrossShopIDConflictDoesNotOverwriteCommittedShop(t *testing.T) {
db := syncTestDB(t)
client := testClient(t, func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/shops") {
writeSyncShops(w, "a", "b", "c")
return
}
shop := r.URL.Query().Get("platformShopId")
id := "a-product"
if shop == "c" {
id = "c-product"
}
writeResponse(w, 200, syncTestPage(shop, 1, 1, id), "")
})
result, err := client.SyncProductData(context.Background(), db, "", nil)
if err != nil || result.SucceededShops != 2 || result.FailedShops != 1 || result.ProductCount != 2 {
t.Fatalf("incorrect conflict result: %+v %v", result, err)
}
assertCode(t, result.Failures[0].Error, "INVALID_RESPONSE")
p, _, _ := db.GetProduct("a-product")
if p.PlatformShopID != "a" {
t.Fatal("conflicting ID overwritten")
}
assertSyncProduct(t, db, false)
}
func TestSyncLocalTransactionFailureContinuesNextShop(t *testing.T) {
db := syncTestDB(t)
if _, err := db.DB().Exec(`CREATE TRIGGER fictional_diagnosis_failure BEFORE INSERT ON product_diagnoses WHEN NEW.product_id='a-product' BEGIN SELECT RAISE(ABORT, 'fictional failure'); END`); err != nil {
t.Fatal(err)
}
client := testClient(t, func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/shops") {
writeSyncShops(w, "a", "b")
return
}
shop := r.URL.Query().Get("platformShopId")
writeResponse(w, 200, syncTestPage(shop, 1, 1, shop+"-product"), "")
})
result, err := client.SyncProductData(context.Background(), db, "", nil)
if err != nil || result.SucceededShops != 1 || result.FailedShops != 1 || result.ProductCount != 1 || result.StopError != nil {
t.Fatalf("local failure lost later shop: %+v %v", result, err)
}
assertCode(t, result.Failures[0].Error, "LOCAL_SYNC_FAILED")
assertSyncProduct(t, db, true)
if _, found, err := db.GetProduct("b-product"); err != nil || !found {
t.Fatal("later shop not committed")
}
}