mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-29 19:36:36 +08:00
232 lines
7.2 KiB
Go
232 lines
7.2 KiB
Go
package service
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
|
|
"go.uber.org/zap"
|
|
|
|
"github.com/ShukeBta/MediaStationGo/internal/model"
|
|
)
|
|
|
|
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")
|
|
}
|
|
|
|
cfg, err := s.repo.StorageConfig.Get(ctx, mount.Provider)
|
|
if err != nil || cfg == nil {
|
|
return res, fmt.Errorf("storage config not found: %s", mount.Provider)
|
|
}
|
|
if !cfg.Enabled {
|
|
return res, fmt.Errorf("storage %s is disabled", mount.Provider)
|
|
}
|
|
typ := mount.Provider
|
|
rootDir := mount.ScanDir
|
|
rootDisplayDir := mount.DisplayDir
|
|
autoCategoryRoot := cloudRootMountNeedsAutoCategory(mount)
|
|
scopeIDs := s.cloudScanLibraryScopeIDs(ctx, lib, mount)
|
|
seen := make(map[string]struct{})
|
|
seenRefs := make(map[string]struct{})
|
|
candidates := make([]cloudCandidate, 0, 256)
|
|
candidateByKey := make(map[string]int)
|
|
visitedDirs := map[string]struct{}{}
|
|
progress := newCloudScanProgressState()
|
|
var stateMu sync.Mutex
|
|
progress.publish(s, lib.ID, res, "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 {
|
|
progress.addSkipped(res)
|
|
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
|
|
}
|
|
setWalkErr(err)
|
|
return err
|
|
}
|
|
progress.publish(s, lib.ID, res, "listing", progress.markDirVisited())
|
|
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) != "" {
|
|
walkWG.Add(1)
|
|
go func(childID, childDisplay string, childMeta *LocalMetadata) {
|
|
_ = walkCloud(childID, childDisplay, childMeta)
|
|
}(entry.ID, joinCloudDisplayPath(displayDir, entry.Name), dirMeta)
|
|
}
|
|
continue
|
|
}
|
|
ext := strings.ToLower(filepath.Ext(entry.Name))
|
|
if _, ok := videoExtensions[ext]; !ok {
|
|
continue
|
|
}
|
|
ref := cloudEntryRef(typ, entry.ID, entry.PickCode)
|
|
if ref == "" {
|
|
progress.addSkipped(res)
|
|
continue
|
|
}
|
|
stateMu.Lock()
|
|
if _, ok := seenRefs[ref]; ok {
|
|
stateMu.Unlock()
|
|
progress.addSkipped(res)
|
|
continue
|
|
}
|
|
seenRefs[ref] = struct{}{}
|
|
stateMu.Unlock()
|
|
progress.publish(s, lib.ID, res, "listing", progress.markFileDiscovered())
|
|
displayPath := joinCloudDisplayPath(displayDir, entry.Name)
|
|
path := cloudMediaPath(typ, displayPath)
|
|
localMeta := s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib))
|
|
localMeta = s.enrichCloudMetadataFromExternalIDs(ctx, lib, path, localMeta)
|
|
if localMeta != nil {
|
|
s.cacheCloudMetadataArtworkNow(ctx, localMeta)
|
|
}
|
|
candidate := cloudCandidate{
|
|
ref: ref,
|
|
name: entry.Name,
|
|
size: entry.Size,
|
|
path: path,
|
|
localMeta: localMeta,
|
|
}
|
|
if autoCategoryRoot {
|
|
candidate.categoryDisplayDir = cloudAutoCategoryDisplayDirForMediaPath(path)
|
|
}
|
|
key := cloudMediaDedupeKey(lib, displayDir, entry.Name, entry.Size)
|
|
stateMu.Lock()
|
|
if key != "" {
|
|
if prevIndex, ok := candidateByKey[key]; ok {
|
|
if candidate.size > candidates[prevIndex].size {
|
|
candidates[prevIndex] = candidate
|
|
}
|
|
stateMu.Unlock()
|
|
progress.addSkipped(res)
|
|
continue
|
|
}
|
|
candidateByKey[key] = len(candidates)
|
|
}
|
|
candidates = append(candidates, candidate)
|
|
stateMu.Unlock()
|
|
}
|
|
return 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.existingCloudMediaSnapshotForLibraries(ctx, scopeIDs)
|
|
if err != nil {
|
|
s.log.Warn("load existing cloud media snapshot failed", zap.String("library_id", lib.ID), zap.Error(err))
|
|
existingMedia = nil
|
|
}
|
|
sortCloudCandidatesByRefreshPriority(candidates, existingMedia)
|
|
writeBatch := newLocalMediaWriteBatch(s, ctx, res, 100)
|
|
probeBudget := maxCloudMediaProbeQueuePerScan
|
|
targetLibs := map[string]*model.Library{"": lib}
|
|
touchedLibraryIDs := []string{}
|
|
for _, candidate := range candidates {
|
|
select {
|
|
case <-ctx.Done():
|
|
return res, ctx.Err()
|
|
default:
|
|
}
|
|
targetLib := lib
|
|
if candidate.categoryDisplayDir != "" {
|
|
if cached, ok := targetLibs[candidate.categoryDisplayDir]; ok {
|
|
targetLib = cached
|
|
} else if categoryLib, err := s.ensureCloudAutoCategoryLibrary(ctx, lib, typ, candidate.categoryDisplayDir); err == nil && categoryLib != nil {
|
|
targetLib = categoryLib
|
|
targetLibs[candidate.categoryDisplayDir] = categoryLib
|
|
scopeIDs = appendUniqueLibraryIDs(scopeIDs, categoryLib.ID)
|
|
} else if err != nil {
|
|
s.log.Warn("ensure cloud auto category library failed",
|
|
zap.String("library_id", lib.ID),
|
|
zap.String("provider", typ),
|
|
zap.String("category", candidate.categoryDisplayDir),
|
|
zap.Error(err))
|
|
}
|
|
}
|
|
touchedLibraryIDs = appendUniqueLibraryIDs(touchedLibraryIDs, targetLib.ID)
|
|
seen[candidate.path] = struct{}{}
|
|
s.ingestCloudFile(ctx, targetLib, typ, candidate.ref, candidate.path, candidate.name, candidate.size, candidate.localMeta, existingMedia, writeBatch, &probeBudget, res)
|
|
progress.publish(s, lib.ID, res, "importing", res.Visited == 1 || res.Visited%100 == 0)
|
|
}
|
|
writeBatch.Flush()
|
|
removed, err := s.pruneMissingCloudMediaForLibraries(ctx, scopeIDs, seen)
|
|
if err != nil {
|
|
s.log.Warn("prune missing cloud media failed", zap.String("library_id", lib.ID), zap.Error(err))
|
|
} else {
|
|
res.Removed = removed
|
|
}
|
|
publishCloudScanFinished(s, lib.ID, res, progress)
|
|
s.invalidateMediaCache(ctx)
|
|
for _, targetID := range appendUniqueLibraryIDs(touchedLibraryIDs, lib.ID) {
|
|
s.maybeGenerateSTRMAfterScan(targetID)
|
|
}
|
|
if scanHasImportChanges(res) && autoScrape && s.scraper != nil && s.scraper.AnyEnabled() && s.autoScrapeEnabled(ctx) {
|
|
for _, targetID := range appendUniqueLibraryIDs(touchedLibraryIDs, lib.ID) {
|
|
s.startAutoScrape(ctx, targetID)
|
|
}
|
|
}
|
|
return res, nil
|
|
}
|
|
|
|
func scanHasImportChanges(res *ScanResult) bool {
|
|
return res != nil && (res.Added > 0 || res.Updated > 0 || res.Removed > 0)
|
|
}
|