fix(subscription): only enqueue missing media

This commit is contained in:
ShukeBta
2026-06-07 13:36:49 +08:00
parent d77e890f51
commit 1d88a09568
3 changed files with 343 additions and 27 deletions
+4
View File
@@ -19,6 +19,7 @@ type LocalAvailability struct {
LocalMediaCount int
MissingEpisodes []int
InLibrary bool
HasSeriesPack bool
ExistingEpisodeKeys map[string]struct{}
MissingEpisodeKeys map[string]struct{}
}
@@ -88,6 +89,9 @@ func LookupLocalAvailability(ctx context.Context, repo *repository.Container, ti
seriesLike := isSubscriptionSeriesType(mediaType)
for _, row := range rows {
if row.EpisodeNum <= 0 {
if seriesLike {
out.HasSeriesPack = true
}
continue
}
season := row.SeasonNum
+123 -27
View File
@@ -205,11 +205,17 @@ func (s *SubscriptionService) runOne(ctx context.Context, sub *model.Subscriptio
seenSet[g] = struct{}{}
}
// 非洗版订阅:成功下载一次即满足,预先算一次媒体库可用性用于跳过已入库的电影/剧集(对齐 MoviePilot)。
// 非洗版订阅:成功下载一次即满足,预先算一次媒体库与下载中任务的
// 可用性,用于跳过已入库/已在下载队列中的电影或剧集(对齐 MoviePilot)。
washOff := !sub.WashEnabled
var avail LocalAvailability
availQuery := ""
if washOff {
avail = SubscriptionLocalAvailability(ctx, s.repo, sub)
availQuery = availabilityQuery(subscriptionName(sub), subscriptionFilter(sub))
avail = mergeLocalAvailability(
SubscriptionLocalAvailability(ctx, s.repo, sub),
s.pendingDownloadAvailability(ctx, sub),
)
}
queued := 0
@@ -226,8 +232,6 @@ func (s *SubscriptionService) runOne(ctx context.Context, sub *model.Subscriptio
continue
}
if washOff && subscriptionItemAlreadyAvailable(sub, avail, item.Title) {
seen = append(seen, guid)
seenSet[guid] = struct{}{}
continue
}
download := item.Enclosure.URL
@@ -240,8 +244,9 @@ func (s *SubscriptionService) runOne(ctx context.Context, sub *model.Subscriptio
mediaType, mediaCategory := s.classifySubscriptionItem(ctx, sub, item.Title, "")
savePath := s.resolveSubscriptionSavePath(ctx, sub, mediaType, mediaCategory)
if s.downloadPathHasCandidate(ctx, sub, item.Title, savePath) {
seen = append(seen, guid)
seenSet[guid] = struct{}{}
if washOff {
addAvailabilityTitle(item.Title, availQuery, &avail)
}
continue
}
if _, err := s.downloads.AddDownloadWithMeta(ctx, sub.UserID, download, savePath, DownloadTaskMeta{
@@ -251,6 +256,9 @@ func (s *SubscriptionService) runOne(ctx context.Context, sub *model.Subscriptio
Overview: sub.Overview,
}); err != nil {
if errors.Is(err, ErrDownloadAlreadyExists) {
if washOff {
addAvailabilityTitle(item.Title, availQuery, &avail)
}
seen = append(seen, guid)
seenSet[guid] = struct{}{}
continue
@@ -263,6 +271,9 @@ func (s *SubscriptionService) runOne(ctx context.Context, sub *model.Subscriptio
zap.Error(err))
continue
}
if washOff {
addAvailabilityTitle(item.Title, availQuery, &avail)
}
queued++
seen = append(seen, guid)
seenSet[guid] = struct{}{}
@@ -413,18 +424,20 @@ func selectSiteSearchCandidates(results []SearchResult, sub *model.Subscription,
Score: score,
})
}
if len(candidates) <= 1 {
return candidates
if len(candidates) > 1 {
sort.SliceStable(candidates, func(i, j int) bool {
if candidates[i].Score != candidates[j].Score {
return candidates[i].Score > candidates[j].Score
}
if candidates[i].Item.Seeders != candidates[j].Item.Seeders {
return candidates[i].Item.Seeders > candidates[j].Item.Seeders
}
return candidates[i].Item.Size > candidates[j].Item.Size
})
}
if len(candidates) == 0 {
return nil
}
sort.SliceStable(candidates, func(i, j int) bool {
if candidates[i].Score != candidates[j].Score {
return candidates[i].Score > candidates[j].Score
}
if candidates[i].Item.Seeders != candidates[j].Item.Seeders {
return candidates[i].Item.Seeders > candidates[j].Item.Seeders
}
return candidates[i].Item.Size > candidates[j].Item.Size
})
var local LocalAvailability
if len(availability) > 0 {
@@ -440,6 +453,9 @@ func selectSiteSearchCandidates(results []SearchResult, sub *model.Subscription,
return candidates[:1]
}
if local.HasSeriesPack {
return nil
}
if local.LocalMediaCount > 0 {
if local.TotalEpisodes > 0 && len(local.MissingEpisodes) == 0 {
return nil
@@ -809,19 +825,95 @@ func (s *SubscriptionService) pendingDownloadAvailability(ctx context.Context, s
if sub != nil {
out.TotalEpisodes = sub.TotalEpisodes
}
root := s.subscriptionBaseSavePath(ctx, sub)
query := availabilityQuery(subscriptionName(sub), subscriptionFilter(sub))
if root == "" || query == "" {
return out
if query == "" {
return s.finalizePendingAvailability(sub, out)
}
_ = scanDownloadPath(ctx, root, query, func(_ string, season, episode int) bool {
out.LocalMediaCount++
out.InLibrary = true
if episode > 0 {
out.ExistingEpisodeKeys[episodeKey(season, episode)] = struct{}{}
root := s.subscriptionBaseSavePath(ctx, sub)
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{}{}
}
return true
})
}
s.addDownloadTaskAvailability(ctx, sub, query, &out)
s.addLiveTorrentAvailability(ctx, query, &out)
return s.finalizePendingAvailability(sub, out)
}
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
}
rows, err := s.repo.Download.List(ctx)
if err != nil {
return
}
baseSavePath := s.subscriptionBaseSavePath(ctx, sub)
for _, row := range rows {
if !downloadTaskBlocksReadd(row.Status) {
continue
}
if baseSavePath != "" && row.SavePath != "" && !sameOrChildPath(row.SavePath, baseSavePath) && !sameOrChildPath(baseSavePath, row.SavePath) {
continue
}
addAvailabilityTitle(row.Title, query, out)
}
}
func (s *SubscriptionService) addLiveTorrentAvailability(ctx context.Context, query string, out *LocalAvailability) {
if s == nil || s.downloads == nil || s.downloads.qb == nil || out == nil {
return
}
live, err := s.downloads.qb.List(ctx, "")
if err != nil {
return
}
for _, torrent := range live {
addAvailabilityTitle(torrent.Name, query, out)
}
}
func addAvailabilityTitle(title, query string, out *LocalAvailability) {
if out == nil || strings.TrimSpace(title) == "" || strings.TrimSpace(query) == "" {
return
}
if !strings.Contains(normalizeAvailabilityComparable(title), normalizeAvailabilityComparable(query)) {
return
}
out.LocalMediaCount++
out.InLibrary = true
season, episode := ParseEpisode(title)
if episode > 0 {
out.ExistingEpisodeKeys[episodeKey(season, episode)] = struct{}{}
return
}
if isSeriesPackTitle(title) {
out.HasSeriesPack = true
}
}
func sameOrChildPath(pathValue, root string) bool {
pathValue = filepath.Clean(strings.TrimSpace(pathValue))
root = filepath.Clean(strings.TrimSpace(root))
if pathValue == "" || root == "" || pathValue == "." || root == "." {
return false
}
if strings.EqualFold(pathValue, root) {
return true
})
}
rel, err := filepath.Rel(root, pathValue)
if err != nil {
return false
}
return rel != "." && !strings.HasPrefix(rel, "..") && !filepath.IsAbs(rel)
}
func (s *SubscriptionService) finalizePendingAvailability(sub *model.Subscription, out LocalAvailability) LocalAvailability {
mediaType := ""
if sub != nil {
mediaType = sub.MediaType
@@ -877,6 +969,7 @@ func mergeLocalAvailability(values ...LocalAvailability) LocalAvailability {
}
out.LocalMediaCount += value.LocalMediaCount
out.InLibrary = out.InLibrary || value.InLibrary
out.HasSeriesPack = out.HasSeriesPack || value.HasSeriesPack
for key := range value.ExistingEpisodeKeys {
out.ExistingEpisodeKeys[key] = struct{}{}
}
@@ -900,12 +993,15 @@ func mergeLocalAvailability(values ...LocalAvailability) LocalAvailability {
// subscriptionItemAlreadyAvailable 判断某个订阅条目(按其标题解析出的季/集)是否已在媒体库存在。
// 电影/无集号条目:媒体库已有该片即视为已存在;剧集条目:对应季集已入库即视为已存在。
func subscriptionItemAlreadyAvailable(sub *model.Subscription, avail LocalAvailability, title string) bool {
if avail.LocalMediaCount == 0 {
if avail.LocalMediaCount == 0 && !avail.HasSeriesPack {
return false
}
if !isSubscriptionSeriesType(subscriptionMediaType(sub)) {
return true
}
if avail.HasSeriesPack {
return true
}
wantSeason, wantEpisode := ParseEpisode(title)
if wantEpisode <= 0 {
// 整季合集 / 无法解析集号:库里已有内容时保守跳过,避免重复整季下载。
+216
View File
@@ -167,6 +167,55 @@ func TestSelectSiteSearchCandidatesWithUnknownTotalSkipsExistingEpisodes(t *test
}
}
func TestSelectSiteSearchCandidatesSingleExistingEpisodeIsSkipped(t *testing.T) {
sub := &model.Subscription{Name: "葬送的芙莉莲 自动订阅", Filter: "葬送的芙莉莲", MediaType: "anime", TotalEpisodes: 3}
results := []SearchResult{
{Title: "葬送的芙莉莲 S01E01 1080p", DownloadURL: "https://pt/download/1", Seeders: 90},
}
availability := LocalAvailability{
TotalEpisodes: 3,
LocalMediaCount: 1,
MissingEpisodes: []int{2, 3},
ExistingEpisodeKeys: map[string]struct{}{episodeKey(1, 1): {}},
}
got := selectSiteSearchCandidates(results, sub, map[string]struct{}{}, availability)
if len(got) != 0 {
t.Fatalf("selected %#v, want none because E01 already exists", got)
}
}
func TestSelectSiteSearchCandidatesSinglePackIsSkippedWhenLibraryPartiallyExists(t *testing.T) {
sub := &model.Subscription{Name: "间谍过家家 自动订阅", Filter: "间谍过家家", MediaType: "tv", TotalEpisodes: 3}
results := []SearchResult{
{Title: "间谍过家家 S01 Complete 1080p", DownloadURL: "https://pt/download/pack", Seeders: 100},
}
availability := LocalAvailability{
TotalEpisodes: 3,
LocalMediaCount: 2,
MissingEpisodes: []int{3},
ExistingEpisodeKeys: map[string]struct{}{episodeKey(1, 1): {}, episodeKey(1, 2): {}},
}
got := selectSiteSearchCandidates(results, sub, map[string]struct{}{}, availability)
if len(got) != 0 {
t.Fatalf("selected %#v, want none because a full pack would redownload existing episodes", got)
}
}
func TestSelectSiteSearchCandidatesSingleExistingMovieIsSkippedWhenNotWashing(t *testing.T) {
sub := &model.Subscription{Name: "Inception 自动订阅", Filter: "Inception 2010", MediaType: "movie"}
results := []SearchResult{
{Title: "Inception 2010 1080p WEB-DL", DownloadURL: "https://pt/download/web", Seeders: 90},
}
availability := LocalAvailability{LocalMediaCount: 1, InLibrary: true, DownloadedEpisodes: 1, TotalEpisodes: 1}
got := selectSiteSearchCandidates(results, sub, map[string]struct{}{}, availability)
if len(got) != 0 {
t.Fatalf("selected %#v, want none because movie already exists and wash is disabled", got)
}
}
func TestSubscriptionPendingDownloadAvailabilitySkipsUnorganizedEpisodes(t *testing.T) {
root := t.TempDir()
seasonDir := filepath.Join(root, "间谍过家家", "Season 01")
@@ -219,6 +268,92 @@ func TestSubscriptionPendingDownloadAvailabilitySkipsUnorganizedEpisodes(t *test
}
}
func TestSubscriptionPendingDownloadAvailabilityIncludesQueuedTasks(t *testing.T) {
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := db.AutoMigrate(&model.DownloadTask{}); 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:2222222222222222222222222222222222222222",
Title: "间谍过家家 S01E02 1080p",
SavePath: "/downloads/tv",
Status: "queued",
}); err != nil {
t.Fatal(err)
}
svc := NewSubscriptionService(nil, nil, repos, nil, nil, nil)
sub := &model.Subscription{
Name: "间谍过家家 自动订阅",
Filter: "间谍过家家",
MediaType: "tv",
SavePath: "/downloads/tv",
TotalEpisodes: 3,
}
availability := svc.pendingDownloadAvailability(t.Context(), sub)
if availability.DownloadedEpisodes != 1 {
t.Fatalf("downloaded episodes = %d, want 1", availability.DownloadedEpisodes)
}
if _, ok := availability.ExistingEpisodeKeys[episodeKey(1, 2)]; !ok {
t.Fatalf("missing queued E02 key: %#v", availability.ExistingEpisodeKeys)
}
results := []SearchResult{
{Title: "间谍过家家 S01E02 1080p WEB-DL", DownloadURL: "https://pt/download/2", Seeders: 80},
{Title: "间谍过家家 S01E03 1080p WEB-DL", DownloadURL: "https://pt/download/3", Seeders: 70},
}
got := selectSiteSearchCandidates(results, sub, map[string]struct{}{}, availability)
if len(got) != 1 || got[0].Episode != 3 {
t.Fatalf("selected %#v, want only not-yet-downloaded episode 3", got)
}
}
func TestSubscriptionPendingDownloadAvailabilityIncludesLiveQBTorrents(t *testing.T) {
qb := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api/v2/auth/login":
_, _ = w.Write([]byte("Ok."))
case "/api/v2/torrents/info":
_, _ = w.Write([]byte(`[{"hash":"abc123","name":"间谍过家家 S01E01 1080p","state":"downloading","progress":0.2}]`))
default:
http.NotFound(w, r)
}
}))
defer qb.Close()
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := db.AutoMigrate(&model.DownloadTask{}); err != nil {
t.Fatal(err)
}
repos := repository.New(db)
downloads := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil)
downloads.qb.Configure(QBitConfig{BaseURL: qb.URL, Username: "admin", Password: "admin"})
svc := NewSubscriptionService(nil, nil, repos, downloads, nil, nil)
sub := &model.Subscription{
Name: "间谍过家家 自动订阅",
Filter: "间谍过家家",
MediaType: "tv",
SavePath: "/downloads/tv",
TotalEpisodes: 2,
}
availability := svc.pendingDownloadAvailability(t.Context(), sub)
if availability.DownloadedEpisodes != 1 {
t.Fatalf("downloaded episodes = %d, want 1", availability.DownloadedEpisodes)
}
if _, ok := availability.ExistingEpisodeKeys[episodeKey(1, 1)]; !ok {
t.Fatalf("missing live qB E01 key: %#v", availability.ExistingEpisodeKeys)
}
}
func TestSubscriptionRunOneDeduplicatesDuplicateRSSGUIDInSameFeed(t *testing.T) {
rss := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/rss+xml")
@@ -299,6 +434,87 @@ func TestSubscriptionRunOneDeduplicatesDuplicateRSSGUIDInSameFeed(t *testing.T)
}
}
func TestSubscriptionRunOneSkipsSameEpisodeAddedEarlierInFeed(t *testing.T) {
rss := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/rss+xml")
_, _ = w.Write([]byte(`<?xml version="1.0"?>
<rss><channel>
<item>
<title>Some Show S01E01 1080p</title>
<guid>episode-1-a</guid>
<link>magnet:?xt=urn:btih:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa&amp;dn=Some+Show+S01E01+1080p</link>
</item>
<item>
<title>Some Show S01E01 WEB-DL</title>
<guid>episode-1-b</guid>
<link>magnet:?xt=urn:btih:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb&amp;dn=Some+Show+S01E01+WEB-DL</link>
</item>
</channel></rss>`))
}))
defer rss.Close()
var addCalls int32
qb := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api/v2/auth/login":
_, _ = w.Write([]byte("Ok."))
case "/api/v2/torrents/info":
if atomic.LoadInt32(&addCalls) > 0 {
_, _ = w.Write([]byte(`[{"hash":"abc123","name":"Some Show S01E01 1080p","state":"downloading","progress":0.1}]`))
return
}
_, _ = w.Write([]byte(`[]`))
case "/api/v2/torrents/add":
atomic.AddInt32(&addCalls, 1)
_, _ = w.Write([]byte("Ok."))
default:
http.NotFound(w, r)
}
}))
defer qb.Close()
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := db.AutoMigrate(&model.Subscription{}, &model.Setting{}, &model.DownloadTask{}, &model.Media{}); err != nil {
t.Fatal(err)
}
repos := repository.New(db)
downloads := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil)
downloads.qb.Configure(QBitConfig{BaseURL: qb.URL, Username: "admin", Password: "admin"})
svc := NewSubscriptionService(nil, zap.NewNop(), repos, downloads, nil, NewHub(zap.NewNop()))
sub := &model.Subscription{
Name: "Some Show 自动订阅",
FeedURL: rss.URL,
Filter: "Some Show",
MediaType: "tv",
SavePath: "/downloads/tv",
TotalEpisodes: 12,
}
if err := repos.Subscription.Create(t.Context(), sub); err != nil {
t.Fatal(err)
}
queued, err := svc.runOne(t.Context(), sub)
if err != nil {
t.Fatal(err)
}
if queued != 1 {
t.Fatalf("queued = %d, want 1", queued)
}
if got := atomic.LoadInt32(&addCalls); got != 1 {
t.Fatalf("qb add calls = %d, want 1", got)
}
rows, err := repos.Download.List(t.Context())
if err != nil {
t.Fatal(err)
}
if len(rows) != 1 {
t.Fatalf("download rows = %d, want 1", len(rows))
}
}
func TestMatchesSubscriptionRulesUserExcludeWords(t *testing.T) {
sub := &model.Subscription{ExcludeWords: "10bit,dolby vision,杜比"}
cases := []struct {