mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-05 13:06:36 +08:00
refactor: split cloud scan candidate collection
This commit is contained in:
@@ -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
|
||||||
|
})
|
||||||
|
}
|
||||||
@@ -3,9 +3,6 @@ package service
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"path/filepath"
|
|
||||||
"strings"
|
|
||||||
"sync"
|
|
||||||
|
|
||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
|
|
||||||
@@ -31,142 +28,17 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
|
|||||||
autoCategoryRoot := cloudRootMountNeedsAutoCategory(mount)
|
autoCategoryRoot := cloudRootMountNeedsAutoCategory(mount)
|
||||||
scopeIDs := s.cloudScanLibraryScopeIDs(ctx, lib, mount)
|
scopeIDs := s.cloudScanLibraryScopeIDs(ctx, lib, mount)
|
||||||
seen := make(map[string]struct{})
|
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()
|
progress := newCloudScanProgressState()
|
||||||
var stateMu sync.Mutex
|
|
||||||
progress.publish(s, lib.ID, res, "listing", true)
|
progress.publish(s, lib.ID, res, "listing", true)
|
||||||
var walkWG sync.WaitGroup
|
candidates, err := s.collectCloudScanCandidates(ctx, lib, cloudScanCandidateRequest{
|
||||||
var walkErr error
|
provider: typ,
|
||||||
var walkErrOnce sync.Once
|
rootDir: rootDir,
|
||||||
setWalkErr := func(err error) {
|
rootDisplayDir: rootDisplayDir,
|
||||||
if err != nil {
|
autoCategoryRoot: autoCategoryRoot,
|
||||||
walkErrOnce.Do(func() {
|
progress: progress,
|
||||||
walkErr = err
|
result: res,
|
||||||
})
|
})
|
||||||
}
|
if err != nil {
|
||||||
}
|
|
||||||
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
|
return res, err
|
||||||
}
|
}
|
||||||
existingMedia, err := s.existingCloudMediaSnapshotForLibraries(ctx, scopeIDs)
|
existingMedia, err := s.existingCloudMediaSnapshotForLibraries(ctx, scopeIDs)
|
||||||
|
|||||||
Reference in New Issue
Block a user