From 829b88036ae80ee17024c16f8ddecd2fc2f25977 Mon Sep 17 00:00:00 2001 From: ShukeBta <272197458+ShukeBta@users.noreply.github.com> Date: Tue, 16 Jun 2026 15:21:19 +0800 Subject: [PATCH] Parallelize cloud library listing --- internal/config/config.go | 20 +++-- internal/service/runtime_settings.go | 10 +++ internal/service/scanner.go | 112 +++++++++++++++++++++---- internal/service/scanner_cloud_test.go | 91 ++++++++++++++++++++ web/src/pages/LibraryPage.tsx | 60 +++++++++++++ 5 files changed, 272 insertions(+), 21 deletions(-) diff --git a/internal/config/config.go b/internal/config/config.go index ed5f97c..49d0718 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -80,11 +80,14 @@ type AppConfig struct { // FFprobeMaxConcurrent limits concurrent ffprobe/ffmpeg metadata probes. // NAS devices can become unresponsive when a scan starts many probe // processes at once, so the default is deliberately conservative. - FFprobeMaxConcurrent int `mapstructure:"ffprobe_max_concurrent"` - MaxCPUThreads int `mapstructure:"max_cpu_threads"` - VAAPIDevice string `mapstructure:"vaapi_device"` - CORSOrigins []string `mapstructure:"cors_origins"` - ServerURL string `mapstructure:"server_url"` + FFprobeMaxConcurrent int `mapstructure:"ffprobe_max_concurrent"` + // CloudScanMaxConcurrent limits concurrent cloud directory list requests + // inside one mounted cloud library scan. + CloudScanMaxConcurrent int `mapstructure:"cloud_scan_max_concurrent"` + MaxCPUThreads int `mapstructure:"max_cpu_threads"` + VAAPIDevice string `mapstructure:"vaapi_device"` + CORSOrigins []string `mapstructure:"cors_origins"` + ServerURL string `mapstructure:"server_url"` } // DatabaseConfig 配置 GORM 数据库。默认 auto: @@ -239,6 +242,7 @@ func setDefaults(v *viper.Viper) { v.SetDefault("app.ffmpeg_path", "ffmpeg") v.SetDefault("app.ffprobe_path", "ffprobe") v.SetDefault("app.ffprobe_max_concurrent", 1) + v.SetDefault("app.cloud_scan_max_concurrent", 4) v.SetDefault("app.max_cpu_threads", 2) v.SetDefault("app.vaapi_device", "/dev/dri/renderD128") v.SetDefault("app.cors_origins", []string{}) @@ -346,6 +350,12 @@ func (c *Config) normalize() error { if c.App.MaxCPUThreads > 8 { c.App.MaxCPUThreads = 8 } + if c.App.CloudScanMaxConcurrent < 1 { + c.App.CloudScanMaxConcurrent = 1 + } + if c.App.CloudScanMaxConcurrent > 16 { + c.App.CloudScanMaxConcurrent = 16 + } if c.Database.MaxOpenConns <= 0 { c.Database.MaxOpenConns = defaultDatabaseMaxOpenConns } diff --git a/internal/service/runtime_settings.go b/internal/service/runtime_settings.go index edffca0..7f83d5e 100644 --- a/internal/service/runtime_settings.go +++ b/internal/service/runtime_settings.go @@ -49,6 +49,16 @@ func ApplyRuntimeSetting(cfg *config.Config, key, value string) { } cfg.App.FFprobeMaxConcurrent = n } + case "cloud.scan_max_concurrent", "cloud.scan_list_concurrency", "app.cloud_scan_max_concurrent": + if n, err := strconv.Atoi(value); err == nil { + if n < 1 { + n = 1 + } + if n > 16 { + n = 16 + } + cfg.App.CloudScanMaxConcurrent = n + } case "app.max_cpu_threads", "runtime.max_cpu_threads": if n, err := strconv.Atoi(value); err == nil { if n < 1 { diff --git a/internal/service/scanner.go b/internal/service/scanner.go index 7bbe3ae..cac6df8 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -1044,50 +1044,90 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar lastProgress := time.Time{} dirsVisited := 0 filesDiscovered := 0 + var stateMu sync.Mutex publishProgress := func(stage string, force bool) { if s.hub == nil { return } + stateMu.Lock() if !force && time.Since(lastProgress) < 2*time.Second { + stateMu.Unlock() return } lastProgress = time.Now() + dirs := dirsVisited + discovered := filesDiscovered + visited := res.Visited + added := res.Added + updated := res.Updated + skipped := res.Skipped + removed := res.Removed + stateMu.Unlock() elapsed := time.Since(startedAt) filesPerSecond := 0.0 - processed := filesDiscovered - if res.Visited > processed { - processed = res.Visited + processed := discovered + if visited > processed { + processed = visited } if elapsed.Seconds() > 0 { filesPerSecond = float64(processed) / elapsed.Seconds() } - s.updateCloudScanProgress(lib.ID, stage, dirsVisited, filesDiscovered, res.Visited, res.Added, res.Updated, res.Skipped, res.Removed, filesPerSecond) + s.updateCloudScanProgress(lib.ID, stage, dirs, discovered, visited, added, updated, skipped, removed, filesPerSecond) s.hub.Publish("scan", map[string]any{ "library_id": lib.ID, "cloud": true, "stage": stage, - "dirs": dirsVisited, - "discovered": filesDiscovered, - "visited": res.Visited, - "added": res.Added, - "updated": res.Updated, - "skipped": res.Skipped, + "dirs": dirs, + "discovered": discovered, + "visited": visited, + "added": added, + "updated": updated, + "skipped": skipped, "elapsed_seconds": int(elapsed.Seconds()), "files_per_second": filesPerSecond, "estimate_message": "云盘接口不提供总文件数,剩余时间会随目录大小和网盘响应速度变化", }) } publishProgress("listing", true) + var walkWG sync.WaitGroup + var walkErr error + var walkErrOnce sync.Once + setWalkErr := func(err error) { + if err != nil { + walkErrOnce.Do(func() { + walkErr = err + }) + } + } + listSlots := make(chan struct{}, s.cloudScanWorkerCount()) var walkCloud func(dirID, displayDir string, inheritedMeta *LocalMetadata) error walkCloud = func(dirID, displayDir string, inheritedMeta *LocalMetadata) error { + defer walkWG.Done() + if err := ctx.Err(); err != nil { + setWalkErr(err) + return err + } + stateMu.Lock() if _, ok := visitedDirs[dirID]; ok { + stateMu.Unlock() return nil } visitedDirs[dirID] = struct{}{} + stateMu.Unlock() + + select { + case listSlots <- struct{}{}: + defer func() { <-listSlots }() + case <-ctx.Done(): + setWalkErr(ctx.Err()) + return ctx.Err() + } entries, err := s.storage.CloudList(ctx, typ, dirID) if err != nil { if dirID != rootDir { + stateMu.Lock() res.Skipped++ + stateMu.Unlock() s.log.Warn("skip inaccessible cloud directory", zap.String("library_id", lib.ID), zap.String("provider", typ), @@ -1095,24 +1135,30 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar zap.Error(err)) return nil } + setWalkErr(err) return err } + stateMu.Lock() dirsVisited++ - publishProgress("listing", dirsVisited == 1 || dirsVisited%20 == 0) + dirProgress := dirsVisited == 1 || dirsVisited%20 == 0 + stateMu.Unlock() + publishProgress("listing", dirProgress) sidecars := newCloudSidecarSet(typ, entries) dirMeta := s.cloudDirectoryMetadata(ctx, typ, displayDir, sidecars, inheritedMeta) s.cacheCloudMetadataArtworkNow(ctx, dirMeta) for _, entry := range entries { select { case <-ctx.Done(): + setWalkErr(ctx.Err()) return ctx.Err() default: } if entry.IsDir { if strings.TrimSpace(entry.ID) != "" { - if err := walkCloud(entry.ID, joinCloudDisplayPath(displayDir, entry.Name), dirMeta); err != nil { - return err - } + walkWG.Add(1) + go func(childID, childDisplay string, childMeta *LocalMetadata) { + _ = walkCloud(childID, childDisplay, childMeta) + }(entry.ID, joinCloudDisplayPath(displayDir, entry.Name), dirMeta) } continue } @@ -1122,16 +1168,22 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar } ref := cloudEntryRef(typ, entry.ID, entry.PickCode) if ref == "" { + stateMu.Lock() res.Skipped++ + stateMu.Unlock() continue } + stateMu.Lock() if _, ok := seenRefs[ref]; ok { res.Skipped++ + stateMu.Unlock() continue } seenRefs[ref] = struct{}{} filesDiscovered++ - publishProgress("listing", filesDiscovered%100 == 0) + fileProgress := filesDiscovered%100 == 0 + stateMu.Unlock() + publishProgress("listing", fileProgress) displayPath := joinCloudDisplayPath(displayDir, entry.Name) path := cloudMediaPath(typ, displayPath) localMeta := s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib)) @@ -1150,21 +1202,32 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar localMeta: localMeta, } key := cloudMediaDedupeKey(lib, displayDir, entry.Name, entry.Size) + stateMu.Lock() if key != "" { if prevIndex, ok := candidateByKey[key]; ok { res.Skipped++ if candidate.size > candidates[prevIndex].size { candidates[prevIndex] = candidate } + stateMu.Unlock() continue } candidateByKey[key] = len(candidates) } candidates = append(candidates, candidate) + stateMu.Unlock() } return nil } - if err := walkCloud(rootDir, rootDisplayDir, nil); err != nil { + walkWG.Add(1) + go func() { + _ = walkCloud(rootDir, rootDisplayDir, nil) + }() + walkWG.Wait() + if walkErr != nil { + return res, walkErr + } + if err := ctx.Err(); err != nil { return res, err } existingMedia, err := s.existingCloudMediaSnapshot(ctx, lib.ID) @@ -1553,6 +1616,23 @@ func (s *ScannerService) ffprobeWorkerCount() int { return normalizeFFprobeMaxConcurrent(s.cfg.App.FFprobeMaxConcurrent) } +func (s *ScannerService) cloudScanWorkerCount() int { + if s == nil || s.cfg == nil { + return 4 + } + return normalizeCloudScanMaxConcurrent(s.cfg.App.CloudScanMaxConcurrent) +} + +func normalizeCloudScanMaxConcurrent(n int) int { + if n <= 0 { + return 1 + } + if n > 16 { + return 16 + } + return n +} + func (s *ScannerService) probeCloudFileMetadata(ctx context.Context, typ, ref string) (*ProbeResult, error) { if s == nil || s.probe == nil || s.storage == nil { return nil, errors.New("cloud probe unavailable") diff --git a/internal/service/scanner_cloud_test.go b/internal/service/scanner_cloud_test.go index 0570569..d4b95db 100644 --- a/internal/service/scanner_cloud_test.go +++ b/internal/service/scanner_cloud_test.go @@ -1,12 +1,16 @@ package service import ( + "context" "encoding/json" "fmt" "net/http" "net/http/httptest" "strings" + "sync" + "sync/atomic" "testing" + "time" "github.com/glebarez/sqlite" "go.uber.org/zap" @@ -122,6 +126,93 @@ func TestScanCloudLibraryImportsRecursivePlayableMedia(t *testing.T) { } } +func TestScanCloudLibraryListsChildDirectoriesConcurrently(t *testing.T) { + var active int32 + var maxActive int32 + var releaseOnce sync.Once + release := make(chan struct{}) + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/file/sort" { + t.Errorf("unexpected path %s", r.URL.Path) + return + } + w.Header().Set("Content-Type", "application/json") + switch r.URL.Query().Get("pdir_fid") { + case "0": + _, _ = w.Write([]byte(`{"status":200,"code":0,"data":{"list":[ + {"fid":"d1","file_name":"A","dir":true,"size":0}, + {"fid":"d2","file_name":"B","dir":true,"size":0} + ]}}`)) + case "d1", "d2": + cur := atomic.AddInt32(&active, 1) + defer atomic.AddInt32(&active, -1) + for { + prev := atomic.LoadInt32(&maxActive) + if cur <= prev || atomic.CompareAndSwapInt32(&maxActive, prev, cur) { + break + } + } + if cur >= 2 { + releaseOnce.Do(func() { close(release) }) + } + select { + case <-release: + case <-r.Context().Done(): + return + case <-time.After(1500 * time.Millisecond): + t.Errorf("child directory requests were not concurrent") + return + } + id := r.URL.Query().Get("pdir_fid") + _, _ = fmt.Fprintf(w, `{"status":200,"code":0,"data":{"list":[{"fid":"f-%s","file_name":"Movie.%s.mkv","dir":false,"size":123}]}}`, id, id) + default: + t.Errorf("unexpected pdir_fid %q", r.URL.Query().Get("pdir_fid")) + } + })) + defer upstream.Close() + + db, err := gorm.Open(sqlite.Open("file:cloud_scan_concurrent?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.Library{}, &model.Media{}, &model.Setting{}, &model.StorageConfig{}); err != nil { + t.Fatal(err) + } + repos := repository.New(db) + log := zap.NewNop() + storage := NewStorageConfigService(log, repos, NewCryptoService("", log)) + if _, err := storage.Save(t.Context(), StorageInput{ + Type: "quark", + Config: map[string]any{ + "cookie": "kps=test", + "base": upstream.URL, + }, + }); err != nil { + t.Fatal(err) + } + lib := model.Library{Name: "夸克网盘", Path: "cloud://quark/0", Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + cfg := &config.Config{} + cfg.App.CloudScanMaxConcurrent = 2 + scanner := NewScannerService(cfg, log, repos, NewHub(log), nil, nil) + scanner.SetStorageConfig(storage) + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Second) + defer cancel() + + res, err := scanner.ScanLibrary(ctx, lib.ID) + if err != nil { + t.Fatalf("scan cloud: %v", err) + } + if got := atomic.LoadInt32(&maxActive); got < 2 { + t.Fatalf("max concurrent child lists = %d, want >= 2", got) + } + if res.Visited != 2 || res.Added != 2 { + t.Fatalf("scan result = %#v, want visited=2 added=2", res) + } +} + func TestCloudLibraryPathParsing(t *testing.T) { typ, dir, ok := parseCloudLibraryPath("cloud://cloud115/abc%20123?ignored=1") if !ok || typ != "cloud115" || dir != "abc 123" { diff --git a/web/src/pages/LibraryPage.tsx b/web/src/pages/LibraryPage.tsx index 65b358b..aa929cd 100644 --- a/web/src/pages/LibraryPage.tsx +++ b/web/src/pages/LibraryPage.tsx @@ -5,6 +5,7 @@ import toast from 'react-hot-toast' import { ArrowLeft, Play, Film } from 'lucide-react' import { libraryAPI } from '../api/library' +import { storageAPI, type CloudScanStatus } from '../api/storage_config' import type { Library, Media } from '../types' import { MediaCard } from '../components/MediaCard' import { ExternalPlayerButton } from '../components/ExternalPlayerButton' @@ -147,6 +148,51 @@ export function LibraryPage() { useWebSocket(onRealtimeEvent) + useEffect(() => { + if (role !== 'admin' || !id) return + let cancelled = false + let terminal = false + let timer: number | undefined + const restoreCloudScanStatus = async () => { + const r = await storageAPI.cloudScanStatus() + if (cancelled) return + const status = (r.items ?? []).find((item) => item.library_id === id) + if (!status) return + if (status.state === 'running' || status.state === 'queued' || status.state === 'canceling') { + setScanning(true) + setScanProgress(formatCloudScanStatus(status)) + return + } + if (status.state === 'finished') { + setScanning(false) + setScanProgress(formatCloudScanStatus(status)) + reloadCurrentLibrary() + terminal = true + if (timer) window.clearInterval(timer) + return + } + if (status.state === 'error' && status.error) { + setScanning(false) + setScanProgress(`扫描失败:${status.error}`) + terminal = true + if (timer) window.clearInterval(timer) + } + } + restoreCloudScanStatus() + .catch(() => undefined) + .finally(() => { + if (!cancelled && !terminal) { + timer = window.setInterval(() => { + restoreCloudScanStatus().catch(() => undefined) + }, 5000) + } + }) + return () => { + cancelled = true + if (timer) window.clearInterval(timer) + } + }, [id, reloadCurrentLibrary, role]) + useEffect(() => { if (loading) return if (!isSeries) { @@ -419,6 +465,20 @@ function formatDuration(seconds: number): string { return `${hours}小时${minutes % 60}分` } +function formatCloudScanStatus(status: CloudScanStatus): string { + const stage = + status.state === 'queued' ? '扫描已排队' + : status.state === 'canceling' ? '正在中断扫描' + : status.state === 'finished' ? '扫描完成' + : status.stage === 'importing' ? '正在入库' + : '正在遍历目录' + const speed = Number(status.files_per_second ?? 0) + const speedText = speed > 0 && status.state !== 'finished' + ? ` · ${speed.toFixed(speed >= 10 ? 0 : 1)} 个/秒` + : '' + return `${stage}:目录 ${status.dirs ?? 0} · 已发现 ${status.discovered ?? 0} · 已入库 ${status.visited ?? 0} · 新增 ${status.added ?? 0} · 更新 ${status.updated ?? 0}${speedText}` +} + function formatSize(bytes: number): string { if (!bytes || bytes <= 0) return '—' const units = ['B', 'KB', 'MB', 'GB', 'TB']