mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-07 05:46:38 +08:00
fix: improve cloud mount scanning feedback
This commit is contained in:
+237
-38
@@ -17,6 +17,7 @@ import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
|
||||
@@ -93,13 +94,35 @@ func (s *ScannerService) ScanLibrary(ctx context.Context, libraryID string) (*Sc
|
||||
return s.scanLibrary(ctx, libraryID, true)
|
||||
}
|
||||
|
||||
// ScanLibraryWithoutAutoScrape walks a library without kicking off online
|
||||
// metadata enrichment. Cloud mounts can contain very large trees; keeping mount
|
||||
// scans import-only prevents scraper bursts from overwhelming small NAS boxes.
|
||||
func (s *ScannerService) ScanLibraryWithoutAutoScrape(ctx context.Context, libraryID string) (*ScanResult, error) {
|
||||
return s.scanLibrary(ctx, libraryID, false)
|
||||
}
|
||||
|
||||
func (s *ScannerService) scanLibrary(ctx context.Context, libraryID string, autoScrape bool) (*ScanResult, error) {
|
||||
lib, err := s.repo.Library.FindByID(ctx, libraryID)
|
||||
if err != nil || lib == nil {
|
||||
return nil, err
|
||||
}
|
||||
if typ, dirID, ok := parseCloudLibraryPath(lib.Path); ok {
|
||||
return s.scanCloudLibrary(ctx, lib, typ, dirID, autoScrape)
|
||||
if mount, ok := ParseCloudLibraryMount(lib.Path); ok {
|
||||
if shadow := s.shadowedCloudLibrary(ctx, lib); shadow != nil {
|
||||
res := &ScanResult{LibraryID: lib.ID, Skipped: 1}
|
||||
s.log.Warn("skip shadowed cloud library scan",
|
||||
zap.String("library_id", lib.ID),
|
||||
zap.String("shadowed_by", shadow.Library.ID),
|
||||
zap.String("provider", mount.Provider))
|
||||
s.hub.Publish("scan", map[string]any{
|
||||
"library_id": lib.ID,
|
||||
"finished": true,
|
||||
"skipped": res.Skipped,
|
||||
"cloud": true,
|
||||
"shadowed": true,
|
||||
})
|
||||
return res, nil
|
||||
}
|
||||
return s.scanCloudLibrary(ctx, lib, mount, autoScrape)
|
||||
}
|
||||
res := &ScanResult{LibraryID: lib.ID}
|
||||
seen := make(map[string]struct{})
|
||||
@@ -174,23 +197,86 @@ func (s *ScannerService) IngestPath(ctx context.Context, libraryID, path string)
|
||||
return res.Added+res.Updated > 0, nil
|
||||
}
|
||||
|
||||
func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Library, typ, rootDir string, autoScrape bool) (*ScanResult, error) {
|
||||
func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Library, mount CloudMountInfo, autoScrape bool) (*ScanResult, error) {
|
||||
res := &ScanResult{LibraryID: lib.ID}
|
||||
if s.storage == nil {
|
||||
return res, fmt.Errorf("cloud storage service unavailable")
|
||||
}
|
||||
typ := mount.Provider
|
||||
rootDir := mount.ScanDir
|
||||
rootDisplayDir := mount.DisplayDir
|
||||
type cloudCandidate struct {
|
||||
ref string
|
||||
name string
|
||||
size int64
|
||||
path string
|
||||
localMeta *LocalMetadata
|
||||
}
|
||||
seen := make(map[string]struct{})
|
||||
seenRefs := make(map[string]struct{})
|
||||
candidates := make([]cloudCandidate, 0, 256)
|
||||
candidateByKey := make(map[string]int)
|
||||
visitedDirs := map[string]struct{}{}
|
||||
var walkCloud func(string) error
|
||||
walkCloud = func(dirID string) error {
|
||||
startedAt := time.Now()
|
||||
lastProgress := time.Time{}
|
||||
dirsVisited := 0
|
||||
filesDiscovered := 0
|
||||
publishProgress := func(stage string, force bool) {
|
||||
if s.hub == nil {
|
||||
return
|
||||
}
|
||||
if !force && time.Since(lastProgress) < 2*time.Second {
|
||||
return
|
||||
}
|
||||
lastProgress = time.Now()
|
||||
elapsed := time.Since(startedAt)
|
||||
filesPerSecond := 0.0
|
||||
processed := filesDiscovered
|
||||
if res.Visited > processed {
|
||||
processed = res.Visited
|
||||
}
|
||||
if elapsed.Seconds() > 0 {
|
||||
filesPerSecond = float64(processed) / elapsed.Seconds()
|
||||
}
|
||||
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,
|
||||
"elapsed_seconds": int(elapsed.Seconds()),
|
||||
"files_per_second": filesPerSecond,
|
||||
"estimate_message": "云盘接口不提供总文件数,剩余时间会随目录大小和网盘响应速度变化",
|
||||
})
|
||||
}
|
||||
publishProgress("listing", true)
|
||||
var walkCloud func(dirID, displayDir string, inheritedMeta *LocalMetadata) error
|
||||
walkCloud = func(dirID, displayDir string, inheritedMeta *LocalMetadata) error {
|
||||
if _, ok := visitedDirs[dirID]; ok {
|
||||
return nil
|
||||
}
|
||||
visitedDirs[dirID] = struct{}{}
|
||||
entries, err := s.storage.CloudList(ctx, typ, dirID)
|
||||
if err != nil {
|
||||
if dirID != rootDir {
|
||||
res.Skipped++
|
||||
s.log.Warn("skip inaccessible cloud directory",
|
||||
zap.String("library_id", lib.ID),
|
||||
zap.String("provider", typ),
|
||||
zap.String("dir", dirID),
|
||||
zap.Error(err))
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
dirsVisited++
|
||||
publishProgress("listing", dirsVisited == 1 || dirsVisited%20 == 0)
|
||||
sidecars := newCloudSidecarSet(typ, entries)
|
||||
dirMeta := s.cloudDirectoryMetadata(ctx, typ, displayDir, sidecars, inheritedMeta)
|
||||
for _, entry := range entries {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
@@ -199,7 +285,7 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
|
||||
}
|
||||
if entry.IsDir {
|
||||
if strings.TrimSpace(entry.ID) != "" {
|
||||
if err := walkCloud(entry.ID); err != nil {
|
||||
if err := walkCloud(entry.ID, joinCloudDisplayPath(displayDir, entry.Name), dirMeta); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
@@ -214,15 +300,45 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
|
||||
res.Skipped++
|
||||
continue
|
||||
}
|
||||
path := cloudMediaPath(typ, ref)
|
||||
seen[path] = struct{}{}
|
||||
s.ingestCloudFile(ctx, lib, typ, ref, entry.Name, entry.Size, res)
|
||||
if _, ok := seenRefs[ref]; ok {
|
||||
res.Skipped++
|
||||
continue
|
||||
}
|
||||
seenRefs[ref] = struct{}{}
|
||||
filesDiscovered++
|
||||
publishProgress("listing", filesDiscovered%100 == 0)
|
||||
displayPath := joinCloudDisplayPath(displayDir, entry.Name)
|
||||
path := cloudMediaPath(typ, displayPath)
|
||||
candidate := cloudCandidate{
|
||||
ref: ref,
|
||||
name: entry.Name,
|
||||
size: entry.Size,
|
||||
path: path,
|
||||
localMeta: s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib)),
|
||||
}
|
||||
key := cloudMediaDedupeKey(lib, displayDir, entry.Name, entry.Size)
|
||||
if key != "" {
|
||||
if prevIndex, ok := candidateByKey[key]; ok {
|
||||
res.Skipped++
|
||||
if candidate.size > candidates[prevIndex].size {
|
||||
candidates[prevIndex] = candidate
|
||||
}
|
||||
continue
|
||||
}
|
||||
candidateByKey[key] = len(candidates)
|
||||
}
|
||||
candidates = append(candidates, candidate)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if err := walkCloud(rootDir); err != nil {
|
||||
if err := walkCloud(rootDir, rootDisplayDir, nil); err != nil {
|
||||
return res, err
|
||||
}
|
||||
for _, candidate := range candidates {
|
||||
seen[candidate.path] = struct{}{}
|
||||
s.ingestCloudFile(ctx, lib, typ, candidate.ref, candidate.path, candidate.name, candidate.size, candidate.localMeta, res)
|
||||
publishProgress("importing", res.Visited == 1 || res.Visited%100 == 0)
|
||||
}
|
||||
removed, err := s.pruneMissingCloudMedia(ctx, lib.ID, seen)
|
||||
if err != nil {
|
||||
s.log.Warn("prune missing cloud media failed", zap.String("library_id", lib.ID), zap.Error(err))
|
||||
@@ -230,13 +346,16 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
|
||||
res.Removed = removed
|
||||
}
|
||||
s.hub.Publish("scan", map[string]any{
|
||||
"library_id": lib.ID,
|
||||
"finished": true,
|
||||
"visited": res.Visited,
|
||||
"added": res.Added,
|
||||
"updated": res.Updated,
|
||||
"removed": res.Removed,
|
||||
"cloud": true,
|
||||
"library_id": lib.ID,
|
||||
"finished": true,
|
||||
"visited": res.Visited,
|
||||
"added": res.Added,
|
||||
"updated": res.Updated,
|
||||
"removed": res.Removed,
|
||||
"discovered": filesDiscovered,
|
||||
"dirs": dirsVisited,
|
||||
"elapsed_seconds": int(time.Since(startedAt).Seconds()),
|
||||
"cloud": true,
|
||||
})
|
||||
if autoScrape && s.scraper != nil && s.scraper.AnyEnabled() && s.autoScrapeEnabled(ctx) {
|
||||
go func(libID string) {
|
||||
@@ -248,7 +367,16 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
|
||||
return res, nil
|
||||
}
|
||||
|
||||
func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library, typ, ref, name string, size int64, res *ScanResult) {
|
||||
func (s *ScannerService) shadowedCloudLibrary(ctx context.Context, lib *model.Library) *CloudMountConflict {
|
||||
libs, err := s.repo.Library.List(ctx)
|
||||
if err != nil {
|
||||
s.log.Warn("list libraries for cloud shadow check failed", zap.String("library_id", lib.ID), zap.Error(err))
|
||||
return nil
|
||||
}
|
||||
return CloudLibraryShadowed(libs, *lib)
|
||||
}
|
||||
|
||||
func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library, typ, ref, path, name string, size int64, localMeta *LocalMetadata, res *ScanResult) {
|
||||
res.Visited++
|
||||
ext := strings.ToLower(filepath.Ext(name))
|
||||
title, year := CleanQuery(name)
|
||||
@@ -258,7 +386,15 @@ func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library
|
||||
if title == "" {
|
||||
title = ref
|
||||
}
|
||||
path := cloudMediaPath(typ, ref)
|
||||
parsedSeason, parsedEpisode := ParseEpisode(path)
|
||||
if librarySupportsSeasons(lib) || parsedSeason > 0 || parsedEpisode > 0 {
|
||||
if seriesTitle, seriesYear := cloudSeriesTitleFromMediaPath(path); seriesTitle != "" {
|
||||
title = seriesTitle
|
||||
if seriesYear > 0 {
|
||||
year = seriesYear
|
||||
}
|
||||
}
|
||||
}
|
||||
isNewMedia := !s.mediaPathExists(ctx, path)
|
||||
m := &model.Media{
|
||||
LibraryID: lib.ID,
|
||||
@@ -277,9 +413,12 @@ func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library
|
||||
s.log.Debug("read cloud strm failed", zap.String("ref", ref), zap.Error(err))
|
||||
}
|
||||
}
|
||||
parsedSeason, parsedEpisode := ParseEpisode(name)
|
||||
m.SeasonNum = parsedSeason
|
||||
m.EpisodeNum = parsedEpisode
|
||||
if localMeta != nil {
|
||||
applyLocalMetadata(m, localMeta)
|
||||
res.LocalMetadata++
|
||||
}
|
||||
if err := s.repo.Media.Upsert(ctx, m); err != nil {
|
||||
s.log.Warn("upsert cloud media failed", zap.String("path", path), zap.Error(err))
|
||||
return
|
||||
@@ -299,6 +438,48 @@ func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library
|
||||
})
|
||||
}
|
||||
|
||||
func cloudSeriesTitleFromMediaPath(mediaPath string) (string, int) {
|
||||
displayPath := strings.TrimSpace(mediaPath)
|
||||
if strings.HasPrefix(strings.ToLower(displayPath), "cloud://") {
|
||||
rest := strings.TrimPrefix(displayPath, "cloud://")
|
||||
if idx := strings.Index(rest, "/"); idx >= 0 {
|
||||
displayPath = rest[idx+1:]
|
||||
} else {
|
||||
return "", 0
|
||||
}
|
||||
}
|
||||
displayPath = strings.Trim(strings.ReplaceAll(displayPath, "\\", "/"), "/")
|
||||
if displayPath == "" {
|
||||
return "", 0
|
||||
}
|
||||
parts := strings.Split(displayPath, "/")
|
||||
if len(parts) < 2 {
|
||||
return "", 0
|
||||
}
|
||||
dirs := parts[:len(parts)-1]
|
||||
if len(dirs) == 0 {
|
||||
return "", 0
|
||||
}
|
||||
base := strings.TrimSpace(dirs[len(dirs)-1])
|
||||
usedSeasonFolder := false
|
||||
if seasonFromDir(base) > 0 {
|
||||
usedSeasonFolder = true
|
||||
dirs = dirs[:len(dirs)-1]
|
||||
if len(dirs) == 0 {
|
||||
return "", 0
|
||||
}
|
||||
base = strings.TrimSpace(dirs[len(dirs)-1])
|
||||
}
|
||||
if base == "" || (!usedSeasonFolder && len(dirs) < 2) {
|
||||
return "", 0
|
||||
}
|
||||
title, year := CleanQuery(base)
|
||||
if title == "" {
|
||||
title = base
|
||||
}
|
||||
return strings.TrimSpace(title), year
|
||||
}
|
||||
|
||||
// RemovePath deletes the media row for a path that has disappeared from disk
|
||||
// (incremental delete used by the watcher on Remove/Rename events).
|
||||
func (s *ScannerService) RemovePath(ctx context.Context, path string) (int64, error) {
|
||||
@@ -488,26 +669,11 @@ func (s *ScannerService) pruneMissingCloudMedia(ctx context.Context, libraryID s
|
||||
}
|
||||
|
||||
func parseCloudLibraryPath(raw string) (typ, dirID string, ok bool) {
|
||||
raw = strings.TrimSpace(raw)
|
||||
if !strings.HasPrefix(strings.ToLower(raw), "cloud://") {
|
||||
info, ok := ParseCloudLibraryMount(raw)
|
||||
if !ok {
|
||||
return "", "", false
|
||||
}
|
||||
u, err := url.Parse(raw)
|
||||
if err != nil || strings.ToLower(u.Scheme) != "cloud" {
|
||||
return "", "", false
|
||||
}
|
||||
typ = strings.TrimSpace(u.Host)
|
||||
if typ == "" {
|
||||
return "", "", false
|
||||
}
|
||||
dirID = strings.Trim(strings.TrimSpace(u.Path), "/")
|
||||
if qDir := strings.TrimSpace(u.Query().Get("dir")); qDir != "" {
|
||||
dirID = qDir
|
||||
}
|
||||
if decoded, err := url.PathUnescape(dirID); err == nil {
|
||||
dirID = decoded
|
||||
}
|
||||
return typ, dirID, true
|
||||
return info.Provider, info.ScanDir, true
|
||||
}
|
||||
|
||||
func cloudEntryRef(typ, id, pickCode string) string {
|
||||
@@ -521,6 +687,39 @@ func cloudMediaPath(typ, ref string) string {
|
||||
return "cloud://" + strings.TrimSpace(typ) + "/" + strings.TrimLeft(strings.TrimSpace(ref), "/")
|
||||
}
|
||||
|
||||
func cloudMediaDedupeKey(lib *model.Library, dirID, name string, size int64) string {
|
||||
base := strings.TrimSpace(strings.TrimSuffix(filepath.Base(name), filepath.Ext(name)))
|
||||
if base == "" {
|
||||
return ""
|
||||
}
|
||||
season, episode := ParseEpisode(name)
|
||||
title, year := CleanQuery(name)
|
||||
title = normalizeCloudDedupeText(title)
|
||||
if (season > 0 || episode > 0) && title != "" {
|
||||
return fmt.Sprintf("episode:%s:%s:%d:%d:%d", strings.ToLower(strings.TrimSpace(lib.Type)), title, year, season, episode)
|
||||
}
|
||||
if (season > 0 || episode > 0) && title == "" {
|
||||
return fmt.Sprintf("episode-dir:%s:%s:%d:%d:%d", strings.ToLower(strings.TrimSpace(lib.Type)), normalizeCloudDedupeText(dirID), season, episode, size)
|
||||
}
|
||||
return fmt.Sprintf("file:%s:%d", normalizeCloudDedupeText(base), size)
|
||||
}
|
||||
|
||||
func normalizeCloudDedupeText(value string) string {
|
||||
value = strings.ToLower(strings.TrimSpace(value))
|
||||
if value == "" {
|
||||
return ""
|
||||
}
|
||||
fields := strings.FieldsFunc(value, func(r rune) bool {
|
||||
switch r {
|
||||
case '.', '_', '-', ' ', '\t', '/', '\\', '[', ']', '(', ')':
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
})
|
||||
return strings.Join(fields, " ")
|
||||
}
|
||||
|
||||
func (s *ScannerService) resolveCloudSTRMTarget(ctx context.Context, typ, ref string) (string, error) {
|
||||
if s.storage == nil {
|
||||
return "", nil
|
||||
|
||||
Reference in New Issue
Block a user