diff --git a/internal/service/scanner_cloud_candidates.go b/internal/service/scanner_cloud_candidates.go new file mode 100644 index 0000000..d804888 --- /dev/null +++ b/internal/service/scanner_cloud_candidates.go @@ -0,0 +1,226 @@ +package service + +import ( + "context" + "path/filepath" + "strings" + "sync" + + "go.uber.org/zap" + + "github.com/ShukeBta/MediaStationGo/internal/model" + "github.com/ShukeBta/MediaStationGo/internal/service/cloud" +) + +type cloudScanCandidateRequest struct { + provider string + rootDir string + rootDisplayDir string + autoCategoryRoot bool + progress *cloudScanProgressState + result *ScanResult +} + +func (s *ScannerService) collectCloudScanCandidates(ctx context.Context, lib *model.Library, req cloudScanCandidateRequest) ([]cloudCandidate, error) { + collector := newCloudScanCandidateCollector(s, ctx, lib, req) + return collector.collect() +} + +type cloudScanCandidateCollector struct { + scanner *ScannerService + ctx context.Context + lib *model.Library + req cloudScanCandidateRequest + + mu sync.Mutex + seenRefs map[string]struct{} + visitedDirs map[string]struct{} + candidates []cloudCandidate + candidateByKey map[string]int + + walkWG sync.WaitGroup + walkErr error + walkErrOnce sync.Once + listSlots chan struct{} +} + +func newCloudScanCandidateCollector(s *ScannerService, ctx context.Context, lib *model.Library, req cloudScanCandidateRequest) *cloudScanCandidateCollector { + return &cloudScanCandidateCollector{ + scanner: s, + ctx: ctx, + lib: lib, + req: req, + seenRefs: make(map[string]struct{}), + visitedDirs: map[string]struct{}{}, + candidates: make([]cloudCandidate, 0, 256), + candidateByKey: make(map[string]int), + listSlots: make(chan struct{}, s.cloudScanWorkerCount()), + } +} + +func (c *cloudScanCandidateCollector) collect() ([]cloudCandidate, error) { + c.walkWG.Add(1) + go func() { + _ = c.walk(c.req.rootDir, c.req.rootDisplayDir, nil) + }() + c.walkWG.Wait() + if c.walkErr != nil { + return nil, c.walkErr + } + if err := c.ctx.Err(); err != nil { + return nil, err + } + return c.candidates, nil +} + +func (c *cloudScanCandidateCollector) walk(dirID, displayDir string, inheritedMeta *LocalMetadata) error { + defer c.walkWG.Done() + if err := c.ctx.Err(); err != nil { + c.setWalkErr(err) + return err + } + if !c.markDirectoryVisited(dirID) { + return nil + } + release, err := c.acquireListSlot() + if err != nil { + c.setWalkErr(err) + return err + } + defer release() + + entries, err := c.scanner.storage.CloudList(c.ctx, c.req.provider, dirID) + if err != nil { + return c.handleListError(dirID, err) + } + c.req.progress.publish(c.scanner, c.lib.ID, c.req.result, "listing", c.req.progress.markDirVisited()) + sidecars := newCloudSidecarSet(c.req.provider, entries) + dirMeta := c.scanner.cloudDirectoryMetadata(c.ctx, c.req.provider, displayDir, sidecars, inheritedMeta) + c.scanner.cacheCloudMetadataArtworkNow(c.ctx, dirMeta) + for _, entry := range entries { + if err := c.ctx.Err(); err != nil { + c.setWalkErr(err) + return err + } + if entry.IsDir { + c.queueChildDirectory(displayDir, entry.Name, entry.ID, dirMeta) + continue + } + c.addFileCandidate(displayDir, entry, sidecars, dirMeta) + } + return nil +} + +func (c *cloudScanCandidateCollector) markDirectoryVisited(dirID string) bool { + c.mu.Lock() + defer c.mu.Unlock() + if _, ok := c.visitedDirs[dirID]; ok { + return false + } + c.visitedDirs[dirID] = struct{}{} + return true +} + +func (c *cloudScanCandidateCollector) acquireListSlot() (func(), error) { + select { + case c.listSlots <- struct{}{}: + return func() { <-c.listSlots }, nil + case <-c.ctx.Done(): + return nil, c.ctx.Err() + } +} + +func (c *cloudScanCandidateCollector) handleListError(dirID string, err error) error { + if dirID != c.req.rootDir { + c.req.progress.addSkipped(c.req.result) + c.scanner.log.Warn("skip inaccessible cloud directory", + zap.String("library_id", c.lib.ID), + zap.String("provider", c.req.provider), + zap.String("dir", dirID), + zap.Error(err)) + return nil + } + c.setWalkErr(err) + return err +} + +func (c *cloudScanCandidateCollector) queueChildDirectory(displayDir, entryName, entryID string, dirMeta *LocalMetadata) { + if strings.TrimSpace(entryID) == "" { + return + } + c.walkWG.Add(1) + go func(childID, childDisplay string, childMeta *LocalMetadata) { + _ = c.walk(childID, childDisplay, childMeta) + }(entryID, joinCloudDisplayPath(displayDir, entryName), dirMeta) +} + +func (c *cloudScanCandidateCollector) addFileCandidate(displayDir string, entry cloud.FileEntry, sidecars cloudSidecarSet, dirMeta *LocalMetadata) { + ext := strings.ToLower(filepath.Ext(entry.Name)) + if _, ok := videoExtensions[ext]; !ok { + return + } + ref := cloudEntryRef(c.req.provider, entry.ID, entry.PickCode) + if ref == "" { + c.req.progress.addSkipped(c.req.result) + return + } + if !c.markRefSeen(ref) { + c.req.progress.addSkipped(c.req.result) + return + } + c.req.progress.publish(c.scanner, c.lib.ID, c.req.result, "listing", c.req.progress.markFileDiscovered()) + displayPath := joinCloudDisplayPath(displayDir, entry.Name) + path := cloudMediaPath(c.req.provider, displayPath) + localMeta := c.scanner.cloudFileMetadata(c.ctx, c.req.provider, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(c.lib)) + localMeta = c.scanner.enrichCloudMetadataFromExternalIDs(c.ctx, c.lib, path, localMeta) + if localMeta != nil { + c.scanner.cacheCloudMetadataArtworkNow(c.ctx, localMeta) + } + candidate := cloudCandidate{ + ref: ref, + name: entry.Name, + size: entry.Size, + path: path, + localMeta: localMeta, + } + if c.req.autoCategoryRoot { + candidate.categoryDisplayDir = cloudAutoCategoryDisplayDirForMediaPath(path) + } + c.addCandidate(displayDir, entry, candidate) +} + +func (c *cloudScanCandidateCollector) markRefSeen(ref string) bool { + c.mu.Lock() + defer c.mu.Unlock() + if _, ok := c.seenRefs[ref]; ok { + return false + } + c.seenRefs[ref] = struct{}{} + return true +} + +func (c *cloudScanCandidateCollector) addCandidate(displayDir string, entry cloud.FileEntry, candidate cloudCandidate) { + key := cloudMediaDedupeKey(c.lib, displayDir, entry.Name, entry.Size) + c.mu.Lock() + defer c.mu.Unlock() + if key != "" { + if prevIndex, ok := c.candidateByKey[key]; ok { + if candidate.size > c.candidates[prevIndex].size { + c.candidates[prevIndex] = candidate + } + c.req.progress.addSkipped(c.req.result) + return + } + c.candidateByKey[key] = len(c.candidates) + } + c.candidates = append(c.candidates, candidate) +} + +func (c *cloudScanCandidateCollector) setWalkErr(err error) { + if err == nil { + return + } + c.walkErrOnce.Do(func() { + c.walkErr = err + }) +} diff --git a/internal/service/scanner_cloud_scan.go b/internal/service/scanner_cloud_scan.go index c4bf692..b450cbb 100644 --- a/internal/service/scanner_cloud_scan.go +++ b/internal/service/scanner_cloud_scan.go @@ -3,9 +3,6 @@ package service import ( "context" "fmt" - "path/filepath" - "strings" - "sync" "go.uber.org/zap" @@ -31,142 +28,17 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar 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 { + candidates, err := s.collectCloudScanCandidates(ctx, lib, cloudScanCandidateRequest{ + provider: typ, + rootDir: rootDir, + rootDisplayDir: rootDisplayDir, + autoCategoryRoot: autoCategoryRoot, + progress: progress, + result: res, + }) + if err != nil { return res, err } existingMedia, err := s.existingCloudMediaSnapshotForLibraries(ctx, scopeIDs)