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) }