feat: 配置批量视频上传并发数 (#31)

This commit is contained in:
QiuSW
2026-10-05 18:04:01 +08:00
parent 622473021a
commit 1541ef56d3
9 changed files with 162 additions and 15 deletions
+25 -3
View File
@@ -9,13 +9,14 @@ import (
"fmt"
"io"
"os"
"sync"
"time"
"cmsp/internal/erpgo"
"cmsp/internal/store"
)
// UploadVideos 只在使用者确认后逐件执行;未决操作只能查询,不能生成新键重放写入。
// UploadVideos 只在使用者确认后按配置并行处理不同商品;未决操作只能查询原键。
func (a *App) UploadVideos(productIDs []string) (uploadErr error) {
defer func() {
if uploadErr != nil && a.log != nil {
@@ -37,9 +38,30 @@ func (a *App) UploadVideos(productIDs []string) (uploadErr error) {
if err != nil {
return err
}
for _, id := range ids {
a.uploadOneVideo(a.appContext(), client, id, len(ids))
concurrency := a.cfg.Upload.Concurrency
if concurrency < 1 || concurrency > 4 {
return fmt.Errorf("上传并发数不合法,允许范围 1—4")
}
if concurrency > len(ids) {
concurrency = len(ids)
}
ctx := a.appContext()
jobs := make(chan string)
var workers sync.WaitGroup
for i := 0; i < concurrency; i++ {
workers.Add(1)
go func() {
defer workers.Done()
for id := range jobs {
a.uploadOneVideo(ctx, client, id, len(ids))
}
}()
}
for _, id := range ids {
jobs <- id
}
close(jobs)
workers.Wait()
return nil
}
+81
View File
@@ -8,13 +8,94 @@ import (
"os"
"path/filepath"
"strings"
"sync/atomic"
"testing"
"time"
"cmsp/internal/downloader"
"cmsp/internal/erpgo"
"cmsp/internal/store"
)
func TestUploadBatchRespectsConfiguredConcurrency(t *testing.T) {
for _, limit := range []int{1, 3} {
t.Run(fmt.Sprintf("concurrency-%d", limit), func(t *testing.T) {
var active, maximum, puts atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
parts := strings.Split(r.URL.Path, "/")
if len(parts) < 7 {
t.Errorf("unexpected path: %s", r.URL.Path)
return
}
itemID := parts[5]
switch {
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/video/operation"):
fmt.Fprintf(w, `{"code":200,"data":{"operationId":"op-1","shopeeId":%q,"status":"succeeded","videoUrl":"https://example.invalid/video.mp4"}}`, itemID)
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/video"):
fmt.Fprintf(w, `{"code":200,"data":{"shopeeId":%q,"source":"huohanhan","fetchedAt":"2026-09-30T00:00:00Z","video":[]}}`, itemID)
case r.Method == http.MethodPut && strings.HasSuffix(r.URL.Path, "/video"):
puts.Add(1)
current := active.Add(1)
for old := maximum.Load(); current > old; old = maximum.Load() {
if maximum.CompareAndSwap(old, current) {
break
}
}
time.Sleep(60 * time.Millisecond)
active.Add(-1)
w.WriteHeader(http.StatusAccepted)
fmt.Fprintf(w, `{"code":202,"data":{"operationId":"op-1","shopeeId":%q,"status":"submitted"}}`, itemID)
default:
t.Errorf("unexpected request: %s %s", r.Method, r.URL.Path)
}
}))
defer server.Close()
db, err := store.Open(":memory:")
if err != nil {
t.Fatal(err)
}
defer db.Close()
a := NewApp()
a.db = db
a.cfg.ERPGo.BaseURL, a.cfg.ERPGo.APIKey = server.URL, "fictional-key"
a.cfg.Upload.Concurrency = limit
a.cfg.Download.VideoDir = t.TempDir()
a.probe = func(context.Context, string) (downloader.ProbeResult, error) {
return downloader.ProbeResult{Duration: 20, FormatName: "mp4", Width: 640, Height: 480}, nil
}
var ids []string
for i := 0; i < 5; i++ {
itemID := fmt.Sprintf("%d", 100+i)
internalID := fmt.Sprintf("internal-%d", i)
ids = append(ids, internalID)
if err := db.UpsertProducts([]store.Product{{ID: internalID, ItemID: itemID, UploadStatus: store.UploadPending}}, "2026-09-30 00:00:00"); err != nil {
t.Fatal(err)
}
dir := filepath.Join(a.cfg.Download.VideoDir, itemID)
if err := os.MkdirAll(dir, 0700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(dir, "video.mp4"), []byte("fictional-mp4"), 0600); err != nil {
t.Fatal(err)
}
}
if err := a.UploadVideos(ids); err != nil {
t.Fatal(err)
}
if puts.Load() != 5 || maximum.Load() != int32(limit) {
t.Fatalf("PUT 数量=%d,最大并发=%d,期望并发=%d", puts.Load(), maximum.Load(), limit)
}
for _, id := range ids {
product, _, err := db.GetProduct(id)
if err != nil || product.UploadStatus != store.UploadDone {
t.Fatalf("商品 %s 状态=%s,err=%v", id, product.UploadStatus, err)
}
}
})
}
}
func TestVideoUploadUsesERPGoAndPersistsOutcome(t *testing.T) {
for _, status := range []string{"succeeded", "unknown"} {
t.Run(status, func(t *testing.T) {
+5
View File
@@ -85,3 +85,8 @@ download:
# 网络层下载失败后的重试次数;HTTP 4xx 和 ffprobe 校验失败不会重试。
download_retries: 3
upload:
# 上传并发数:批量上传同时处理的商品数,允许 1—4;默认 1。
# 每个商品仍使用独立幂等键,未决操作只查询原任务,不重复上传。
concurrency: 1
+5 -5
View File
@@ -2,8 +2,8 @@
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: 55761e55811a7eaab87ff41a97de07375425c0ab
synchronized_at: 2026-10-05T02:08:22Z
wiki_revision: 9ebe51001ec9054656478312baedb3387d172c7a
synchronized_at: 2026-10-05T10:01:39Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
@@ -23,13 +23,13 @@ DownloadProductData(platformShopID) 返回 internal/erpgo.SyncResult 和 error
## 本地列表分页条数(2026-09-29)
商品列表 n-pagination 启用 show-size-picker,提供每页20/50/100/200/500条;初始/重置仍20。changePageSize 更新 query.pageSize、清空旧勾选,调用 search 回到第一页,其他筛选条件保留。仅此次页大小切换明确清空勾选,普通搜索/翻页继续原选择行为。
商品列表 n-pagination 启用 show-size-picker,提供每页20/50/100/200/500条;初始/重置仍20。changePageSize 更新 query.pageSize、清空旧勾选,调用 search 回到第一页,其他筛选条件保留。普通搜索/重置和翻页清空勾选;后台刷新只保留当前页仍可见的 ID。上传预览与实际执行使用同一选择快照,不会把隐藏的旧勾选带入批量操作。
store.ListProducts 接受1—500的 pageSize,空/非正数/大于500仍回退20;参数化 LIMIT/OFFSET、排序与同条件 COUNT 保持原样。这里是本地商品列表,与 erpgo 商品同步每页200、任务并发及上传流程无关。选择500增加当前页行和视频摘要读取,不自动勾选或启动任务。
## 商品列表下载状态筛选(2026-09-29)
`ProductListView.vue` 的初始/重置查询默认为在售中(NORMAL)、缺少视频(missing)和未下载(not_downloaded)。下载状态下拉复用现有 n-select,提供全部、未下载(含失败)、待下载、下载中、已下载、下载失败;用户点击搜索应用条件,沿用原页码重置和筛选栏换行;普通搜索/翻页保持原有勾选行为。
`ProductListView.vue` 的初始/重置查询默认为在售中(NORMAL)、缺少视频(missing)和未下载(not_downloaded)。下载状态下拉复用现有 n-select,提供全部、未下载(含失败)、待下载、下载中、已下载、下载失败;用户点击搜索应用条件,沿用原页码重置和筛选栏换行;普通搜索/翻页不会携带不可见行的勾选。
`store.ProductQuery.DownloadStatus` 的 JSON 名称为 `downloadStatus`,空值保持旧调用查全部下载状态。`buildWhere` 把 not_downloaded 映射到 `download_status IN (pending, failed)`,单状态等值查询,所有值参数化,与其他条件 AND 组合;COUNT 和分页列表使用同一 WHERE。not_downloaded 仅是查询条件,不写入数据库状态。筛选不调用外部接口、不更改 SQLite 表结构或商品/视频/任务/上传状态。
@@ -137,7 +137,7 @@ cmsp/
→ 事件推送进度,前端刷新表格
使用者勾选商品点击「上传视频」
→ app_upload.go 串行处理,预检本地 mp4 并查询 erpgo 当前视频
→ app_upload.go 按 `upload.concurrency`(默认 1,允许 1—4)对不同商品有界并行处理;每件仍独立预检本地 mp4 并查询 erpgo 当前视频
→ SQLite 保存幂等键、Shopee ID 与文件指纹
→ internal/erpgo/video.go 以 X-API-Key PUT 单个 multipart video
→ 同键 GET /video/operation 查询结果;仅 succeeded 更新本地视频和上传状态
+6 -4
View File
@@ -2,8 +2,8 @@
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: 337733444b6e824e5c847a6648bbbd70e11ea8c5
synchronized_at: 2026-10-05T02:28:55Z
wiki_revision: bbb42cb37cc4703c3953c021a7054f862ae312e0
synchronized_at: 2026-10-05T10:01:40Z
<!-- gitea-wiki-mirror:end -->
# 业务规则与术语
@@ -12,6 +12,8 @@ synchronized_at: 2026-10-05T02:28:55Z
单商品显示远端已有视频的覆盖警告;批量发现远端已有视频则跳过。上传文件取 Shopee ID 子目录下排序首个 mp4,先完成时长、格式、像素与大小校验。写入前在 SQLite 持久化幂等键。HTTP 202 和 `processing`/`unknown` 均不代表成功;再次操作仅查询原键,`succeeded` 才更新本地视频与上传状态。ERPGo 预检返回 200 且 `video: []` 时才确认远端当前无视频;返回 502 / `HHH_UPSTREAM_ERROR` / `stage=check` 表示视频状态未知,显示“视频状态无法确认,上传未提交”,不创建上传操作、不发送 PUT,保留本地 MP4 与 requestId。批量路径的该商品上传状态为 `unconfirmed`,不记作上传失败。`tempVideoUrl` 或 `videoUploadIdStr` 非空表示远端仍有处理中标记:单商品预览停止,执行前再次检查;批量遇到该商品标为待处理,不创建新操作、不发送 PUT。只有普通已存在视频且没有处理中标记时,单商品才可在明确覆盖确认后继续。用户确认后弹窗立即关闭,上传按钮在调用期间继续显示运行中;弹窗关闭不代表远端成功。单商品已有本地未决操作时,确认框改为“恢复原上传操作”,只查询或恢复原键,不提示覆盖、不发新 PUT;操作明确失败后才按新上传流程检查当前视频。远端操作 processing/unknown 或操作查询异常保留原键,商品上传状态为 `unconfirmed`,列表显示“待确认”;只有本次操作返回 `succeeded` 才显示“已上传”。
批量上传同时处理的不同商品数由本机 `upload.concurrency` 控制,默认 1、允许 1—4;参数设置页可修改并保存到 `config.yaml`。每件商品独立保存原幂等键;未决操作仍只查询、不重新提交。一个批次同一时刻只运行一组任务,用户需先明确勾选并确认。
## 未选择店铺的同步范围(2026-09-29,#26)
@@ -23,7 +25,7 @@ synchronized_at: 2026-10-05T02:28:55Z
## 本地商品列表每页条数(2026-09-29)
分页栏可选20、50、100、200、500条,默认及重置20。切换每页条数后回到第一页、保留店铺/诊断/下载/上传等筛选、清空旧勾选;列表总数仍为全部匹配记录数。500条是显示范围,不表示自动批量处理500条,批量操作仍需勾选。
分页栏可选20、50、100、200、500条,默认及重置20。切换每页条数后回到第一页、保留店铺/诊断/下载/上传等筛选、清空旧勾选;列表总数仍为全部匹配记录数。500条是显示范围,不表示自动批量处理500条,批量操作仍需勾选。新搜索/重置和翻页清空旧勾选;后台刷新移除当前页不可见 ID,表头全选仅覆盖当前页,弹窗和实际处理只使用当前页可见的勾选。
后端本地查询允许1—500条,无效值回退20;空页仍返回匹配总数。分页只读取SQLite,不扫描或改写任务数据,不修改 erpgo 接口分页大小或外部同步规则。
@@ -137,7 +139,7 @@ MTOP 搜索方式、签名与请求参数保持现状。商品原链接以实际
- erpgo 已提供 Shopee 商品视频上传与操作查询契约;实际 Shopee 发布生效时间和独立验证仍待真实业务验收。
- SQLite 表结构与迁移方式。
- 并发下载与上传的默认并发数,以及淘宝风控的实际容忍阈值。
- 淘宝风控的实际容忍阈值。
- 视频与商品的匹配规则:一个商品搜到多个同款时,选哪一个的视频,是否需要人工确认。
- 失败重试策略:哪些失败自动重试、重试几次、哪些必须人工介入。
- Shopee 侧的实际生效延迟与验证方式。
@@ -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: 023ef74aa74316cd47ee50757d356634ddd37c41
synchronized_at: 2026-09-30T01:39:57Z
wiki_revision: ebd03ea3ee41b327b1ed7dd4070f918a586cf4f2
synchronized_at: 2026-10-05T10:01:41Z
<!-- gitea-wiki-mirror:end -->
# 本地开发与验证
## 批量上传并发与选择范围(2026-10-05)
`upload.concurrency` 在本机 `config.yaml` 中配置,默认 1、允许 1—4;`config.example.yaml` 给出示例,参数设置页可修改并保存。旧配置缺字段保持串行;该值与仅控制 CDN 文件下载的 `download.concurrency` 独立。`app_upload.go` 对不同商品有界并行,单商品幂等键、未决操作只查原键、批次互斥及 `succeeded` 才写本地完成状态的规则不变。模拟 HTTP/临时 SQLite 的 `TestUploadBatchRespectsConfiguredConcurrency` 覆盖默认串行和并发 3,配置测试覆盖旧文件默认、保存读回与非法值;真实批量写入及 ERPGo 限流需单独验收。
商品列表搜索、重置、翻页后不沿用隐藏勾选;后台刷新只保留当前页可见 ID。复现“当前筛选 13 条却带入旧选 48 条”时,以列表显示、已选数、确认弹窗及最终调用的同一组 ID 为验收依据。
## erpgo 视频上传验证(2026-09-30)
运行 `go test ./...`、`go test -race ./...`、`go vet ./...`、`npm --prefix frontend run build` 和 Wails 构建。模拟 HTTP 和临时 SQLite 覆盖 multipart `video`、`X-API-Key`、幂等键、结果查询、成功落库及 unknown 不重发;真实店铺写入与 Shopee 页面生效需单独验收。
@@ -21,7 +27,7 @@ internal/erpgo/sync_test.go 使用模拟HTTP与临时SQLite验证最新店铺列
## 每页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。
internal/store/product_test.go 使用501条虚构匹配记录和1条不匹配记录,验证20/50/100/200/500五档有效、无效值回退20、500条末页1条、随后空页、分页不重叠与总数一致。前端使用已安装Naive UI分页控件的showSizePicker/pageSizes/onUpdate:pageSize契约;changePageSize 保留筛选、清空勾选并回到第一页。普通搜索/重置和翻页也清空旧勾选,后台刷新只保留当前页可见 ID;批量预览与执行使用同一快照。默认/重置为20。
执行 go test ./...、go vet ./...、前端构建及 wails build -o cmsp25.exe。构建及数据库测试不代替Windows实际下拉、键盘、Narrator、高对比度/显示缩放或每页500条的真机响应体验验收。
+5
View File
@@ -289,6 +289,11 @@ onUnmounted(() => {
<n-form-item label="视频文件下载并发数(1—8)">
<n-input-number v-model:value="cfg.download.concurrency" :min="1" :max="8" />
</n-form-item>
<n-form-item label="商品视频上传并发数(1—4)">
<n-input-number v-model:value="cfg.upload.concurrency" :min="1" :max="4" />
</n-form-item>
</div>
<div class="two">
<n-form-item label="商品及候选间等待秒数">
<n-input-number v-model:value="cfg.download.waitSecondsMin" :min="0" style="width: 90px" />
<span class="dash">—</span>
+16
View File
@@ -39,6 +39,7 @@ type Config struct {
Huohanhan HuohanhanConfig `yaml:"huohanhan" json:"huohanhan"`
Taobao TaobaoConfig `yaml:"taobao" json:"taobao"`
Download DownloadConfig `yaml:"download" json:"download"`
Upload UploadConfig `yaml:"upload" json:"upload"`
}
// ERPGoConfig 用于查询店铺和商品;APIKey 只保存在本机,禁止进入日志。
@@ -107,6 +108,11 @@ type DownloadConfig struct {
DownloadRetries int `yaml:"download_retries" json:"downloadRetries"`
}
// UploadConfig 控制一次批量视频上传中同时处理的商品数。
type UploadConfig struct {
Concurrency int `yaml:"concurrency" json:"concurrency"`
}
// Default 返回一份可以直接使用的默认配置。
//
// 账号和密码故意留空:程序不内置任何凭据,必须由使用者自己填。
@@ -137,6 +143,7 @@ func Default() Config {
RiskEmptyThreshold: 8,
DownloadRetries: 3,
},
Upload: UploadConfig{Concurrency: 1},
}
}
@@ -238,6 +245,9 @@ func (c Config) Validate() error {
if err := checkIntRange(d.DownloadRetries, 0, 5, "下载重试次数"); err != nil {
return err
}
if err := checkIntRange(c.Upload.Concurrency, 1, 4, "上传并发数"); err != nil {
return err
}
return nil
}
@@ -395,6 +405,11 @@ download:
# 网络层下载失败后的重试次数;HTTP 4xx 和 ffprobe 校验失败不会重试。
download_retries: %d
upload:
# 上传并发数:一次批量上传同时处理的商品数,允许 1—4。
# 每件商品仍使用独立幂等键,未决操作不重复上传。
concurrency: %d
`,
quoted(c.ERPGo.BaseURL),
quoted(c.ERPGo.APIKey),
@@ -416,6 +431,7 @@ download:
trimFloat(c.Download.GuardWaitSeconds),
c.Download.RiskEmptyThreshold,
c.Download.DownloadRetries,
c.Upload.Concurrency,
)
}
+10
View File
@@ -68,6 +68,8 @@ func TestValidateRejectsBadValues(t *testing.T) {
{"每商品视频数为 0", func(c *Config) { c.Download.MaxVideosPerProduct = 0 }, "1—10"},
{"图搜取数过大", func(c *Config) { c.Download.SearchTopN = 61 }, "1—60"},
{"并发数过大", func(c *Config) { c.Download.Concurrency = 9 }, "1—8"},
{"上传并发数为零", func(c *Config) { c.Upload.Concurrency = 0 }, "1—4"},
{"上传并发数过大", func(c *Config) { c.Upload.Concurrency = 5 }, "1—4"},
{"等待区间颠倒", func(c *Config) { c.Download.WaitSecondsMax = 1; c.Download.WaitSecondsMin = 5 }, "不能小于最短"},
}
@@ -116,6 +118,7 @@ func TestSaveThenLoadKeepsValues(t *testing.T) {
saved.Huohanhan.Password = `测试"密码\含转义` // 故意含引号和反斜杠
saved.Download.MaxVideosPerProduct = 5
saved.Download.WaitSecondsMin = 1.5
saved.Upload.Concurrency = 3
if err := Save(path, saved); err != nil {
t.Fatalf("保存失败:%v", err)
@@ -137,6 +140,9 @@ func TestSaveThenLoadKeepsValues(t *testing.T) {
if loaded.Download.WaitSecondsMin != 1.5 {
t.Fatalf("小数应当读回 1.5,实际 %v", loaded.Download.WaitSecondsMin)
}
if loaded.Upload.Concurrency != 3 {
t.Fatalf("上传并发数应当读回 3,实际 %d", loaded.Upload.Concurrency)
}
}
// Windows 路径全是反斜杠,必须能原样存取。
@@ -202,6 +208,7 @@ func TestSaveKeepsComments(t *testing.T) {
"这里没有淘宝账号和密码",
"不会代填密码",
"淘宝页面访问保持串行",
"上传并发数",
} {
if !strings.Contains(text, must) {
t.Fatalf("保存后应当保留注释 %q,实际内容:\n%s", must, text)
@@ -224,6 +231,9 @@ func TestLoadFillsMissingFieldsWithDefaults(t *testing.T) {
if cfg.Download.MaxVideosPerProduct != 7 {
t.Fatalf("已有字段应当读回 7,实际 %d", cfg.Download.MaxVideosPerProduct)
}
if cfg.Upload.Concurrency != 1 {
t.Fatalf("旧配置缺少上传并发数时应保持串行,实际 %d", cfg.Upload.Concurrency)
}
if cfg.Taobao.DebugPortStart != Default().Taobao.DebugPortStart {
t.Fatalf("缺失字段应当保留默认值 %d,实际 %d",
Default().Taobao.DebugPortStart, cfg.Taobao.DebugPortStart)