diff --git a/internal/handler/subscriptions.go b/internal/handler/subscriptions.go index 112b129..db31d34 100644 --- a/internal/handler/subscriptions.go +++ b/internal/handler/subscriptions.go @@ -79,7 +79,7 @@ func createSubscriptionHandler(svc *service.Container) gin.HandlerFunc { return } enriched := []model.Subscription{*s} - service.EnrichSubscriptionProgress(c.Request.Context(), svc.Repo, enriched) + svc.Subscription.EnrichProgress(c.Request.Context(), enriched) *s = enriched[0] c.JSON(http.StatusCreated, s) } @@ -93,7 +93,7 @@ func listSubscriptionsHandler(svc *service.Container) gin.HandlerFunc { return } enrichAndPersistSubscriptions(c.Request.Context(), svc, items) - service.EnrichSubscriptionProgress(c.Request.Context(), svc.Repo, items) + svc.Subscription.EnrichProgress(c.Request.Context(), items) c.JSON(http.StatusOK, gin.H{"items": items}) } } @@ -105,7 +105,7 @@ func listSubscriptionHistoryHandler(svc *service.Container) gin.HandlerFunc { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } - service.EnrichSubscriptionProgress(c.Request.Context(), svc.Repo, items) + svc.Subscription.EnrichProgress(c.Request.Context(), items) c.JSON(http.StatusOK, gin.H{"items": items}) } } diff --git a/internal/service/downloads.go b/internal/service/downloads.go index edda905..c895ffe 100644 --- a/internal/service/downloads.go +++ b/internal/service/downloads.go @@ -490,10 +490,41 @@ func (d *DownloadService) torrentExistsByIdentity(ctx context.Context, title str } func downloadTaskIdentityKey(name string) string { + if key := downloadMediaIdentityKey(name); key != "" { + return key + } + return normalizedDownloadTitleKey(name) +} + +func downloadMediaIdentityKey(name string) string { name = strings.ToLower(strings.TrimSpace(name)) if name == "" { return "" } + title, year := CleanQuery(name) + titleKey := normalizeAvailabilityComparable(title) + if titleKey == "" { + titleKey = normalizeAvailabilityComparable(availabilityQuery(name, "")) + } + if titleKey == "" { + return "" + } + season, episode := ParseEpisode(name) + parts := []string{titleKey} + if year > 0 { + parts = append(parts, fmt.Sprintf("y%d", year)) + } + if episode > 0 { + if season <= 0 { + season = 1 + } + parts = append(parts, fmt.Sprintf("s%02de%03d", season, episode)) + } + return strings.Join(parts, "|") +} + +func normalizedDownloadTitleKey(name string) string { + name = strings.ToLower(strings.TrimSpace(name)) var b strings.Builder for _, r := range name { if unicode.IsLetter(r) || unicode.IsDigit(r) { diff --git a/internal/service/downloads_test.go b/internal/service/downloads_test.go index 0a498d5..e8d9f11 100644 --- a/internal/service/downloads_test.go +++ b/internal/service/downloads_test.go @@ -631,7 +631,7 @@ func TestAddDownloadWithMetaSkipsExistingTaskBeforeQBAdd(t *testing.T) { svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil) svc.qb.Configure(QBitConfig{BaseURL: qb.URL, Username: "admin", Password: "admin"}) task, err := svc.AddDownloadWithMeta(t.Context(), "u1", "https://pt.example/download?id=new&passkey=new", "/downloads/tv", DownloadTaskMeta{ - Title: "Some Show S01E01 1080p", + Title: "Some Show S01E01 2160p WEB-DL", }) if !errors.Is(err, ErrDownloadAlreadyExists) { t.Fatalf("err = %v, want ErrDownloadAlreadyExists", err) diff --git a/internal/service/scanner.go b/internal/service/scanner.go index 77503dc..275edf5 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -365,6 +365,7 @@ type existingLocalMedia struct { AudioCodec string Container string STRMURL string + FileID string } func (s *ScannerService) cloudMediaProbeWorker() { @@ -829,6 +830,12 @@ func (s *ScannerService) scanLibrary(ctx context.Context, libraryID string, auto if err != nil { s.log.Warn("load existing local media snapshot failed", zap.String("library_id", lib.ID), zap.Error(err)) existingMedia = nil + } else { + for path, existing := range existingMedia { + if existing.FileID != "" { + seenInodes[existing.FileID] = path + } + } } walkFn := func(path string, info walkInfo) error { @@ -1124,6 +1131,7 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar return priority(candidates[i]) < priority(candidates[j]) }) } + writeBatch := newLocalMediaWriteBatch(s, ctx, res, 100) probeBudget := maxCloudMediaProbeQueuePerScan for _, candidate := range candidates { select { @@ -1132,9 +1140,10 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar default: } seen[candidate.path] = struct{}{} - s.ingestCloudFile(ctx, lib, typ, candidate.ref, candidate.path, candidate.name, candidate.size, candidate.localMeta, existingMedia, &probeBudget, res) + s.ingestCloudFile(ctx, lib, typ, candidate.ref, candidate.path, candidate.name, candidate.size, candidate.localMeta, existingMedia, writeBatch, &probeBudget, res) publishProgress("importing", res.Visited == 1 || res.Visited%100 == 0) } + writeBatch.Flush() 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)) @@ -1233,10 +1242,11 @@ func (s *ScannerService) existingLocalMediaSnapshot(ctx context.Context, library AudioCodec string Container string STRMURL string + FileID string } if err := s.repo.DB.WithContext(ctx). Model(&model.Media{}). - Select("path, size_bytes, duration_sec, width, height, video_codec, audio_codec, container, strm_url"). + Select("path, size_bytes, duration_sec, width, height, video_codec, audio_codec, container, strm_url, file_id"). Where("library_id = ? AND path NOT LIKE ?", libraryID, "cloud://%"). Find(&rows).Error; err != nil { return nil, err @@ -1253,6 +1263,7 @@ func (s *ScannerService) existingLocalMediaSnapshot(ctx context.Context, library AudioCodec: row.AudioCodec, Container: row.Container, STRMURL: row.STRMURL, + FileID: row.FileID, } } } @@ -1292,7 +1303,7 @@ func (s *ScannerService) shadowedCloudLibrary(ctx context.Context, lib *model.Li return CloudLibraryShadowed(libs, *lib) } -func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library, typ, ref, path, name string, size int64, localMeta *LocalMetadata, existingMedia map[string]existingCloudMedia, probeBudget *int, res *ScanResult) { +func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library, typ, ref, path, name string, size int64, localMeta *LocalMetadata, existingMedia map[string]existingCloudMedia, writeBatch *localMediaWriteBatch, probeBudget *int, res *ScanResult) { res.Visited++ ext := strings.ToLower(filepath.Ext(name)) title, year := CleanQuery(name) @@ -1353,6 +1364,16 @@ func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library s.queueCloudArtworkPrefetch(localMeta.PosterURL) s.queueCloudArtworkPrefetch(localMeta.BackdropURL) } + if isNewMedia && writeBatch != nil { + var after func() + if needsTrackProbe && ext != ".strm" { + after = func() { + s.queueCloudMediaProbeWithBudget(typ, ref, path, probeBudget) + } + } + writeBatch.AddWithAfter(path, m, after) + return + } if err := s.repo.Media.Upsert(ctx, m); err != nil { addScanError(res, path, err) s.log.Warn("upsert cloud media failed", zap.String("path", path), zap.Error(err)) @@ -1485,14 +1506,6 @@ func cloudTrackMetadataMissing(existing existingCloudMedia) bool { strings.TrimSpace(existing.AudioCodec) == "" } -func localTrackMetadataMissing(existing existingLocalMedia) bool { - return existing.DurationSec <= 0 || - existing.Width <= 0 || - existing.Height <= 0 || - strings.TrimSpace(existing.VideoCodec) == "" || - strings.TrimSpace(existing.AudioCodec) == "" -} - func localMetadataNeedsRefresh(local *LocalMetadata) bool { return local != nil && (local.HasNFO || local.HasArtwork || localHasDescriptiveMetadata(local)) } @@ -1559,6 +1572,7 @@ func (s *ScannerService) RemovePath(ctx context.Context, path string) (int64, er func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, path string, size int64, seenInodes map[string]string, existingMedia map[string]existingLocalMedia, writeBatch *localMediaWriteBatch, res *ScanResult) { res.Visited++ ext := strings.ToLower(filepath.Ext(path)) + cleanPath := filepath.Clean(path) // Hardlink dedup: a seeding source kept by keep_seeding shares its inode // with the organized hardlink. Importing both would create duplicate rows @@ -1572,16 +1586,17 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat zap.String("path", path), zap.String("primary", first)) return } - if other, ok := s.duplicateByFileID(ctx, fileID, path); ok { - res.Skipped++ - s.log.Debug("scan skip hardlink duplicate (existing)", - zap.String("path", path), zap.String("primary", other)) - return + if existingMedia == nil { + if other, ok := s.duplicateByFileID(ctx, fileID, path); ok { + res.Skipped++ + s.log.Debug("scan skip hardlink duplicate (existing)", + zap.String("path", path), zap.String("primary", other)) + return + } } seenInodes[fileID] = path } - cleanPath := filepath.Clean(path) parsedSeason, parsedEpisode := ParseEpisode(path) localMeta, localMetaErr := ReadLocalMetadata(path, lib.Path, librarySupportsSeasons(lib) || parsedSeason > 0 || parsedEpisode > 0) if localMetaErr != nil { @@ -1594,7 +1609,6 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat if exists && ext != ".strm" && existing.SizeBytes == size && - (s.probe == nil || !localTrackMetadataMissing(existing)) && !localMetadataNeedsRefresh(localMeta) { res.Skipped++ return @@ -1687,6 +1701,7 @@ type localMediaWriteBatch struct { type localMediaWriteItem struct { path string media *model.Media + after func() } func newLocalMediaWriteBatch(scanner *ScannerService, ctx context.Context, res *ScanResult, limit int) *localMediaWriteBatch { @@ -1697,13 +1712,17 @@ func newLocalMediaWriteBatch(scanner *ScannerService, ctx context.Context, res * } func (b *localMediaWriteBatch) Add(path string, media *model.Media) { + b.AddWithAfter(path, media, nil) +} + +func (b *localMediaWriteBatch) AddWithAfter(path string, media *model.Media, after func()) { if b == nil || b.scanner == nil || media == nil { return } if media.ScrapeStatus == "" { media.ScrapeStatus = "pending" } - b.items = append(b.items, localMediaWriteItem{path: path, media: media}) + b.items = append(b.items, localMediaWriteItem{path: path, media: media, after: after}) if len(b.items) >= b.limit { b.Flush() } @@ -1726,6 +1745,11 @@ func (b *localMediaWriteBatch) Flush() { } if err := b.scanner.repo.DB.WithContext(b.ctx).CreateInBatches(&media, b.limit).Error; err == nil { b.res.Added += len(media) + for _, item := range items { + if item.after != nil { + item.after() + } + } b.publish() return } @@ -1739,6 +1763,9 @@ func (b *localMediaWriteBatch) Flush() { continue } b.res.Added++ + if item.after != nil { + item.after() + } } b.publish() } diff --git a/internal/service/scanner_incremental_test.go b/internal/service/scanner_incremental_test.go index b6d61db..0017aa3 100644 --- a/internal/service/scanner_incremental_test.go +++ b/internal/service/scanner_incremental_test.go @@ -126,6 +126,35 @@ func TestScanLibrarySkipsUnchangedExistingLocalMedia(t *testing.T) { } } +func TestScanLibrarySkipsUnchangedExistingLocalMediaWithMissingTrackMetadata(t *testing.T) { + sc, repos := newScannerTestEnv(t) + root := t.TempDir() + lib := model.Library{Name: "Movies", Path: root, Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + file := filepath.Join(root, "Already In Library Missing Tracks (2024).mkv") + if err := os.WriteFile(file, []byte("same-size"), 0o644); err != nil { + t.Fatal(err) + } + first, err := sc.ScanLibrary(t.Context(), lib.ID) + if err != nil { + t.Fatalf("first scan: %v", err) + } + if first.Added != 1 { + t.Fatalf("first scan = %#v, want added=1", first) + } + + sc.probe = NewFFprobeService(&config.Config{}, zap.NewNop()) + second, err := sc.ScanLibrary(t.Context(), lib.ID) + if err != nil { + t.Fatalf("second scan: %v", err) + } + if second.Added != 0 || second.Updated != 0 || second.Skipped != 1 || second.Probed != 0 { + t.Fatalf("second scan = %#v, want unchanged file skipped without synchronous track probe", second) + } +} + func TestScanLibraryReportsPerFileUpsertErrors(t *testing.T) { sc, repos := newScannerTestEnv(t) root := t.TempDir() diff --git a/internal/service/subscription_availability.go b/internal/service/subscription_availability.go index 989ea47..fa6ba1e 100644 --- a/internal/service/subscription_availability.go +++ b/internal/service/subscription_availability.go @@ -27,7 +27,6 @@ func (s *SubscriptionService) pendingDownloadAvailability(ctx context.Context, s if root != "" { _ = scanDownloadPath(ctx, root, query, func(_ string, season, episode int) bool { out.LocalMediaCount++ - out.InLibrary = true if episode > 0 { out.ExistingEpisodeKeys[episodeKey(season, episode)] = struct{}{} } @@ -39,6 +38,22 @@ func (s *SubscriptionService) pendingDownloadAvailability(ctx context.Context, s return s.finalizePendingAvailability(sub, out) } +func (s *SubscriptionService) EnrichProgress(ctx context.Context, items []model.Subscription) { + for i := range items { + availability := mergeLocalAvailability( + SubscriptionLocalAvailability(ctx, s.repo, &items[i]), + s.pendingDownloadAvailability(ctx, &items[i]), + ) + items[i].DownloadedEpisodes = availability.DownloadedEpisodes + items[i].LocalMediaCount = availability.LocalMediaCount + items[i].MissingEpisodes = availability.MissingEpisodes + items[i].InLibrary = availability.InLibrary + if items[i].TotalEpisodes == 0 { + items[i].TotalEpisodes = availability.TotalEpisodes + } + } +} + func (s *SubscriptionService) addDownloadTaskAvailability(ctx context.Context, sub *model.Subscription, query string, out *LocalAvailability) { if s == nil || s.repo == nil || s.repo.Download == nil || out == nil { return @@ -80,7 +95,6 @@ func addAvailabilityTitle(title, query string, out *LocalAvailability) { return } out.LocalMediaCount++ - out.InLibrary = true season, episode := ParseEpisode(title) if episode > 0 { out.ExistingEpisodeKeys[episodeKey(season, episode)] = struct{}{} diff --git a/internal/service/subscription_test.go b/internal/service/subscription_test.go index d33425a..62aadea 100644 --- a/internal/service/subscription_test.go +++ b/internal/service/subscription_test.go @@ -300,6 +300,9 @@ func TestSubscriptionPendingDownloadAvailabilityIncludesQueuedTasks(t *testing.T if availability.DownloadedEpisodes != 1 { t.Fatalf("downloaded episodes = %d, want 1", availability.DownloadedEpisodes) } + if availability.InLibrary { + t.Fatal("queued download should not be reported as already in library") + } if _, ok := availability.ExistingEpisodeKeys[episodeKey(1, 2)]; !ok { t.Fatalf("missing queued E02 key: %#v", availability.ExistingEpisodeKeys) } @@ -314,6 +317,42 @@ func TestSubscriptionPendingDownloadAvailabilityIncludesQueuedTasks(t *testing.T } } +func TestSubscriptionEnrichProgressIncludesPendingDownloads(t *testing.T) { + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.DownloadTask{}, &model.Media{}); err != nil { + t.Fatal(err) + } + repos := repository.New(db) + if err := repos.Download.Create(t.Context(), &model.DownloadTask{ + Source: "qbittorrent", + URL: "magnet:?xt=urn:btih:4444444444444444444444444444444444444444", + Title: "Inception 2010 1080p", + SavePath: "/downloads/movies", + Status: "completed", + Progress: 1, + }); err != nil { + t.Fatal(err) + } + svc := NewSubscriptionService(nil, nil, repos, nil, nil, nil) + items := []model.Subscription{{ + Name: "Inception 2010", + Filter: "Inception 2010", + MediaType: "movie", + SavePath: "/downloads/movies", + }} + + svc.EnrichProgress(t.Context(), items) + if items[0].InLibrary { + t.Fatal("pending download should not be reported as in-library media") + } + if items[0].DownloadedEpisodes != 1 || items[0].LocalMediaCount != 1 || items[0].TotalEpisodes != 1 { + t.Fatalf("unexpected enriched progress: %+v", items[0]) + } +} + func TestSubscriptionLocalAvailabilityMatchesMediaPath(t *testing.T) { db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) if err != nil { diff --git a/web/src/pages/SubscriptionsPage.tsx b/web/src/pages/SubscriptionsPage.tsx index 4338104..67cb1e2 100644 --- a/web/src/pages/SubscriptionsPage.tsx +++ b/web/src/pages/SubscriptionsPage.tsx @@ -370,7 +370,8 @@ function subscriptionRuleBadges(subscription: Subscription): string[] { function subscriptionProgressLabel(subscription: Subscription): string { const isSeries = ['tv', 'anime', 'variety'].includes((subscription.media_type || '').toLowerCase()) if (!isSeries) { - return subscription.in_library ? '本地已入库' : '本地未入库' + if (subscription.in_library) return '本地已入库' + return (subscription.downloaded_episodes || subscription.local_media_count || 0) > 0 ? '已下载未入库' : '本地未入库' } const downloaded = subscription.downloaded_episodes || 0 const total = subscription.total_episodes || 0