diff --git a/app.go b/app.go index 2562926..b2348f6 100644 --- a/app.go +++ b/app.go @@ -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 启动选中商品的批量取视频任务。 diff --git a/app_live_test.go b/app_live_test.go index 9007d03..7bcee8a 100644 --- a/app_live_test.go +++ b/app_live_test.go @@ -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}) diff --git a/app_sync_test.go b/app_sync_test.go index 4c48178..56fd895 100644 --- a/app_sync_test.go +++ b/app_sync_test.go @@ -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") } diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index 124c296..5f5d32f 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -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 # 架构与代码地图 +## 全部店铺商品同步(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 仍用于现有视频上传,原客户端代码保留用于回归与回退。 ## 当前实现状态 diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 640a8bc..1313756 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -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 # 业务规则与术语 +## 未选择店铺的同步范围(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 签名 | diff --git a/docs/04-local-development-and-verification.md b/docs/04-local-development-and-verification.md index 5ea146e..f64a85f 100644 --- a/docs/04-local-development-and-verification.md +++ b/docs/04-local-development-and-verification.md @@ -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 # 本地开发与验证 +## 全店铺同步验证(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。 diff --git a/docs/06-troubleshooting.md b/docs/06-troubleshooting.md index 4fa2952..eaf997d 100644 --- a/docs/06-troubleshooting.md +++ b/docs/06-troubleshooting.md @@ -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 # 故障排查 +## 全店铺同步与部分成功(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 | 处理 | |---|---| diff --git a/frontend/src/views/ProductListView.vue b/frontend/src/views/ProductListView.vue index 17673f2..dc3e7b7 100644 --- a/frontend/src/views/ProductListView.vue +++ b/frontend/src/views/ProductListView.vue @@ -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(() => { 下载数据 @@ -822,6 +833,30 @@ onUnmounted(() => { + +

+ 成功店铺的数据已保存,失败及未处理店铺的原商品数据和任务状态保留。 +

+

+ 停止原因:{{ queryErrorMessage(JSON.stringify(productSyncResult.stopError)) }} +

+
+ 查看失败店铺({{ productSyncResult.failures.length }}) +
    +
  • + {{ failure.shopName || failure.platformShopId }}:{{ queryErrorMessage(JSON.stringify(failure.error)) }} +
  • +
+
+
+
@@ -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; diff --git a/internal/erpgo/sync.go b/internal/erpgo/sync.go new file mode 100644 index 0000000..cdc38bd --- /dev/null +++ b/internal/erpgo/sync.go @@ -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 + } +} diff --git a/internal/erpgo/sync_test.go b/internal/erpgo/sync_test.go new file mode 100644 index 0000000..013f841 --- /dev/null +++ b/internal/erpgo/sync_test.go @@ -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") + } +}