Parallelize cloud library listing

This commit is contained in:
ShukeBta
2026-06-16 15:21:19 +08:00
parent 18ae6df115
commit 829b88036a
5 changed files with 272 additions and 21 deletions
+15 -5
View File
@@ -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
}
+10
View File
@@ -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 {
+96 -16
View File
@@ -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")
+91
View File
@@ -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" {
+60
View File
@@ -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']