fix subscription stale download dedup

This commit is contained in:
ShukeBta
2026-06-26 17:05:31 +08:00
parent 0e11d64a2f
commit a587054443
8 changed files with 235 additions and 6 deletions
+53
View File
@@ -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)
@@ -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)
}
}
+12 -4
View File
@@ -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 {
-1
View File
@@ -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
}
@@ -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
@@ -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)
}
}
@@ -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)
}
}
@@ -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
}