diff --git a/internal/service/download_add.go b/internal/service/download_add.go index 8077200..7d6ab13 100644 --- a/internal/service/download_add.go +++ b/internal/service/download_add.go @@ -5,6 +5,7 @@ import ( "errors" "path" "strings" + "time" "github.com/ShukeBta/MediaStationGo/internal/model" ) @@ -246,6 +247,9 @@ func (d *DownloadService) findExistingDownloadTask(ctx context.Context, req down if !downloadTaskInSubscriptionScope(rows[i], req) { continue } + if !d.subscriptionDownloadTaskStillLive(ctx, rows[i]) { + continue + } } else if !downloadTaskBlocksDuplicate(rows[i].Status) { continue } @@ -257,6 +261,55 @@ func (d *DownloadService) findExistingDownloadTask(ctx context.Context, req down return nil, false } +func (d *DownloadService) subscriptionDownloadTaskStillLive(ctx context.Context, row model.DownloadTask) bool { + live, ok := d.liveTorrentSnapshot(30 * time.Second) + if !ok && d != nil && d.qb != nil && d.qb.IsConfigured() { + var err error + live, err = d.qb.List(ctx, "") + if err != nil { + return true + } + ok = true + } + if !ok { + return true + } + for _, torrent := range live { + if downloadTaskMatchesLiveTorrent(row, torrent) { + return true + } + } + return false +} + +func downloadTaskMatchesLiveTorrent(row model.DownloadTask, torrent QBitTorrent) bool { + torrentName := strings.TrimSpace(torrent.Name) + if torrentName == "" { + return false + } + req := downloadAddRequest{ + title: row.Title, + savePath: row.SavePath, + meta: DownloadTaskMeta{ + SubscriptionID: row.SubscriptionID, + }, + } + if downloadTaskCoversAddRequest(torrentName, req) { + return true + } + rowKey := downloadTaskIdentityKey(row.Title) + torrentKey := downloadTaskIdentityKey(torrentName) + if rowKey != "" && torrentKey != "" { + return rowKey == torrentKey + } + if len(episodeRefsFromTitle(row.Title)) > 0 || len(episodeRefsFromTitle(torrentName)) > 0 { + return false + } + rowTorrentKey := normalizeTorrentName(row.Title) + liveTorrentKey := normalizeTorrentName(torrentName) + return rowTorrentKey != "" && rowTorrentKey == liveTorrentKey +} + func downloadTaskCoversAddRequest(existing string, req downloadAddRequest) bool { if subscriptionRequestHasExplicitEpisodes(req) { return downloadExplicitEpisodesCoverRequest(existing, req.title) diff --git a/internal/service/download_add_dedup_test.go b/internal/service/download_add_dedup_test.go index 2b650b6..3dd5675 100644 --- a/internal/service/download_add_dedup_test.go +++ b/internal/service/download_add_dedup_test.go @@ -228,3 +228,57 @@ func TestAddDownloadWithMetaScopesSubscriptionDedupBySubscriptionOrSavePath(t *t t.Fatalf("qb add calls = %d, want 1", got) } } + +func TestAddDownloadWithMetaRequeuesStaleSubscriptionTaskMissingFromQB(t *testing.T) { + 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":"newhash","name":"Archives The Nanyang Mystery 2026 S01E07-S01E08 2160p WEB-DL","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 := newServiceTestDB(t, &model.DownloadTask{}, &model.DownloadClient{}, &model.Setting{}) + repos := repository.New(db) + configureTestDefaultQB(t, repos, qb.URL) + if err := repos.Download.Create(t.Context(), &model.DownloadTask{ + UserID: "u1", + SubscriptionID: "sub-nanyang", + Source: "qbittorrent", + URL: "https://pt.example/download?id=stale", + Title: "Archives The Nanyang Mystery 2026 S01E07-S01E08 2160p WEB-DL", + SavePath: "/downloads/tv", + Status: "queued", + Progress: 0, + }); err != nil { + t.Fatal(err) + } + + 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", "magnet:?xt=urn:btih:bcbcbcbcbcbcbcbcbcbcbcbcbcbcbcbcbcbcbcbc&dn=Archives+The+Nanyang+Mystery+2026+S01E07-S01E08", "/downloads/tv", DownloadTaskMeta{ + SubscriptionID: "sub-nanyang", + Title: "Archives The Nanyang Mystery 2026 S01E07-S01E08 2160p WEB-DL", + }) + if err != nil { + t.Fatalf("AddDownloadWithMeta returned %v, want stale task ignored and candidate queued", err) + } + if task == nil { + t.Fatal("task = nil, want requeued task") + } + if got := atomic.LoadInt32(&addCalls); got != 1 { + t.Fatalf("qb add calls = %d, want 1", got) + } +} diff --git a/internal/service/download_live_snapshot.go b/internal/service/download_live_snapshot.go index 4d5334a..bbc67bb 100644 --- a/internal/service/download_live_snapshot.go +++ b/internal/service/download_live_snapshot.go @@ -21,19 +21,27 @@ func (d *DownloadService) recordLiveTorrentSnapshot(live []QBitTorrent) { } func (d *DownloadService) LiveTorrentSnapshot(maxAge time.Duration) []QBitTorrent { - if d == nil { + snapshot, ok := d.liveTorrentSnapshot(maxAge) + if !ok { return nil } + return snapshot +} + +func (d *DownloadService) liveTorrentSnapshot(maxAge time.Duration) ([]QBitTorrent, bool) { + if d == nil { + return nil, false + } now := d.currentTime() d.mu.Lock() defer d.mu.Unlock() if d.liveTorrentsAt.IsZero() { - return nil + return nil, false } if maxAge > 0 && now.Sub(d.liveTorrentsAt) > maxAge { - return nil + return nil, false } - return cloneQBitTorrentSlice(d.liveTorrents) + return cloneQBitTorrentSlice(d.liveTorrents), true } func cloneQBitTorrentSlice(in []QBitTorrent) []QBitTorrent { diff --git a/internal/service/subscription.go b/internal/service/subscription.go index 2b71efe..9acc2ed 100644 --- a/internal/service/subscription.go +++ b/internal/service/subscription.go @@ -315,7 +315,6 @@ func (s *SubscriptionService) enqueueRSSSubscriptionCandidate(ctx context.Contex AllowExistingLibrary: sub.WashEnabled, }); err != nil { if IsDownloadDedupError(err) { - state.markTitleAvailable(item.Title) state.markSeen(candidate.GUID) return false } diff --git a/internal/service/subscription_availability.go b/internal/service/subscription_availability.go index 4cafdc2..542d6a5 100644 --- a/internal/service/subscription_availability.go +++ b/internal/service/subscription_availability.go @@ -71,6 +71,9 @@ func (s *SubscriptionService) addDownloadTaskAvailability(ctx context.Context, s if !downloadTaskBlocksReadd(row.Status) { continue } + if !s.downloadTaskCountsAsPending(ctx, row) { + continue + } linkedToSubscription := sub != nil && strings.TrimSpace(row.SubscriptionID) != "" && row.SubscriptionID == sub.ID if !linkedToSubscription && baseSavePath != "" && row.SavePath != "" && !sameOrChildPath(row.SavePath, baseSavePath) && !sameOrChildPath(baseSavePath, row.SavePath) { continue @@ -83,6 +86,13 @@ func (s *SubscriptionService) addDownloadTaskAvailability(ctx context.Context, s } } +func (s *SubscriptionService) downloadTaskCountsAsPending(ctx context.Context, row model.DownloadTask) bool { + if s == nil || s.downloads == nil { + return true + } + return s.downloads.subscriptionDownloadTaskStillLive(ctx, row) +} + func (s *SubscriptionService) addLiveTorrentAvailability(ctx context.Context, query string, out *LocalAvailability) { if s == nil || s.downloads == nil || s.downloads.qb == nil || out == nil { return diff --git a/internal/service/subscription_availability_test.go b/internal/service/subscription_availability_test.go index f2a358f..5a46173 100644 --- a/internal/service/subscription_availability_test.go +++ b/internal/service/subscription_availability_test.go @@ -225,3 +225,66 @@ func TestSubscriptionPendingDownloadAvailabilityIncludesLiveQBTorrents(t *testin t.Fatalf("missing live qB E01 key: %#v", availability.ExistingEpisodeKeys) } } + +func TestSiteSearchDownloadDedupDoesNotMarkCandidateAvailable(t *testing.T) { + db := newServiceTestDB(t, &model.DownloadTask{}, &model.Setting{}) + repos := repository.New(db) + sub := &model.Subscription{ + Base: model.Base{ID: "sub-nanyang"}, + UserID: "u1", + Name: "南部档案 自动订阅", + Filter: "南部档案", + MediaType: "tv", + SavePath: "/downloads/tv", + TotalEpisodes: 33, + } + if err := repos.Download.Create(t.Context(), &model.DownloadTask{ + SubscriptionID: sub.ID, + Source: "qbittorrent", + URL: "https://pt/download/existing", + Title: "Archives The Nanyang Mystery 2026 S01E07-S01E08 2160p WEB-DL", + SavePath: "/downloads/tv", + Status: "queued", + Progress: 0, + }); err != nil { + t.Fatal(err) + } + downloads := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil) + siteSvc := NewSiteService(zap.NewNop(), repos, "") + svc := NewSubscriptionService(nil, zap.NewNop(), repos, downloads, siteSvc, NewHub(zap.NewNop())) + state := &siteSearchRunState{ + Keyword: "南部档案", + SeenSet: map[string]struct{}{}, + Availability: LocalAvailability{ + TotalEpisodes: 33, + ExistingEpisodeKeys: map[string]struct{}{}, + MissingEpisodeKeys: map[string]struct{}{}, + }, + } + + title, err := svc.enqueueSiteSearchCandidate(t.Context(), sub, siteSearchCandidate{ + Item: SearchResult{ + Title: "Archives The Nanyang Mystery 2026 S01E07-S01E08 2160p WEB-DL", + DownloadURL: "https://pt/download/existing", + }, + Download: "https://pt/download/existing", + GUID: "site|mteam|nanyang-7-8", + Season: 1, + Episode: 7, + Episodes: []int{7, 8}, + Pack: true, + Score: 80, + }, state) + if err != nil { + t.Fatalf("enqueueSiteSearchCandidate returned %v, want dedup skipped without error", err) + } + if title != "" { + t.Fatalf("queued title = %q, want empty on dedup", title) + } + if _, ok := state.Availability.ExistingEpisodeKeys[episodeKey(1, 7)]; ok { + t.Fatalf("deduped candidate should not mark E07 available: %#v", state.Availability.ExistingEpisodeKeys) + } + if len(state.Seen) != 0 { + t.Fatalf("deduped candidate should not be marked seen: %#v", state.Seen) + } +} diff --git a/internal/service/subscription_pending_availability_test.go b/internal/service/subscription_pending_availability_test.go index 22ce3cc..4f6d393 100644 --- a/internal/service/subscription_pending_availability_test.go +++ b/internal/service/subscription_pending_availability_test.go @@ -5,6 +5,8 @@ import ( "path/filepath" "testing" + "go.uber.org/zap" + "github.com/ShukeBta/MediaStationGo/internal/model" "github.com/ShukeBta/MediaStationGo/internal/repository" ) @@ -137,3 +139,44 @@ func TestSubscriptionPendingDownloadAvailabilityIncludesLinkedAliasTask(t *testi t.Fatalf("selected %#v, want linked alias task to satisfy E21", got) } } + +func TestSubscriptionPendingDownloadAvailabilitySkipsStaleTaskMissingFromQB(t *testing.T) { + db := newServiceTestDB(t, &model.DownloadTask{}) + repos := repository.New(db) + sub := &model.Subscription{ + Base: model.Base{ID: "sub-nanyang"}, + Name: "南部档案 自动订阅", + Filter: "南部档案", + MediaType: "tv", + SavePath: "/downloads/tv", + TotalEpisodes: 33, + } + if err := repos.Download.Create(t.Context(), &model.DownloadTask{ + SubscriptionID: sub.ID, + Source: "qbittorrent", + URL: "https://pt/download/stale", + Title: "Archives The Nanyang Mystery 2026 S01E07-S01E08 2160p WEB-DL", + SavePath: "/downloads/tv", + Status: "queued", + Progress: 0, + }); err != nil { + t.Fatal(err) + } + downloads := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil) + downloads.recordLiveTorrentSnapshot(nil) + svc := NewSubscriptionService(nil, nil, repos, downloads, nil, nil) + + availability := svc.pendingDownloadAvailability(t.Context(), sub) + if availability.DownloadedEpisodes != 0 { + t.Fatalf("downloaded episodes = %d, want stale task not counted", availability.DownloadedEpisodes) + } + if _, ok := availability.ExistingEpisodeKeys[episodeKey(1, 7)]; ok { + t.Fatalf("stale E07 task should not count as available: %#v", availability.ExistingEpisodeKeys) + } + got := selectSiteSearchCandidates([]SearchResult{ + {Title: "Archives The Nanyang Mystery 2026 S01E07-S01E08 2160p WEB-DL", SearchKeyword: "南部档案 2026", DownloadURL: "https://pt/download/7-8", Seeders: 80}, + }, sub, map[string]struct{}{}, availability) + if len(got) != 1 || got[0].Episode != 7 { + t.Fatalf("selected %#v, want stale missing range to be eligible", got) + } +} diff --git a/internal/service/subscription_site_search_enqueue.go b/internal/service/subscription_site_search_enqueue.go index 63519dc..6b5579c 100644 --- a/internal/service/subscription_site_search_enqueue.go +++ b/internal/service/subscription_site_search_enqueue.go @@ -68,7 +68,6 @@ func (s *SubscriptionService) enqueueSiteSearchCandidate(ctx context.Context, su AllowExistingLibrary: sub.WashEnabled, }); err != nil { if IsDownloadDedupError(err) { - state.markCandidateAvailable(candidate) s.logSiteSearchCandidateSkipped(sub, state, candidate, "download_dedup", mediaType, mediaCategory, savePath) return "", nil }