From e5b5ecce9387ddf29f8fcd0cfe1a4bf5944475df Mon Sep 17 00:00:00 2001 From: ShukeBta <272197458+ShukeBta@users.noreply.github.com> Date: Wed, 24 Jun 2026 19:48:22 +0800 Subject: [PATCH] refactor: split subscription site search enqueue --- internal/service/subscription_site_search.go | 261 ++++++------------ .../subscription_site_search_enqueue.go | 150 ++++++++++ 2 files changed, 232 insertions(+), 179 deletions(-) create mode 100644 internal/service/subscription_site_search_enqueue.go diff --git a/internal/service/subscription_site_search.go b/internal/service/subscription_site_search.go index ec83044..207e989 100644 --- a/internal/service/subscription_site_search.go +++ b/internal/service/subscription_site_search.go @@ -34,6 +34,65 @@ func (s *SubscriptionService) runSiteSearch(ctx context.Context, sub *model.Subs s.log.Info("site-search subscription run started", subscriptionSiteSearchLogFields(sub, keyword)...) } + results, err := s.searchSubscriptionSites(ctx, sub, keywords) + if err != nil { + return 0, err + } + if len(results) == 0 { + return s.finishSiteSearchNoResults(sub, keyword) + } + s.updateSubscriptionTotalEpisodes(ctx, sub, s.resolveSubscriptionTotalEpisodes(ctx, sub, inferSearchTotalEpisodes(results, sub))) + + guidKey, seen, seenSet := s.loadSiteSearchSeen(ctx, sub) + availability := mergeLocalAvailability( + SubscriptionLocalAvailability(ctx, s.repo, sub), + s.pendingDownloadAvailability(ctx, sub), + ) + candidates, selectionStats := selectSiteSearchCandidatesWithStats(results, sub, seenSet, availability) + if s.log != nil { + fields := subscriptionSiteSearchLogFields(sub, keyword) + fields = appendSiteSearchSelectionLogFields(fields, selectionStats) + fields = appendAvailabilityLogFields(fields, availability) + s.log.Info("site-search subscription selection summary", fields...) + } + runState := &siteSearchRunState{ + Keyword: keyword, + Seen: seen, + SeenSet: seenSet, + Availability: availability, + } + queueResult := s.enqueueSiteSearchCandidates(ctx, sub, candidates, runState) + availability = s.finalizePendingAvailability(sub, runState.Availability) + seen = trimSiteSearchSeen(runState.Seen) + _ = s.repo.Setting.Set(ctx, guidKey, strings.Join(seen, "\n")) + now := time.Now() + _ = s.repo.DB.Model(sub).Updates(map[string]any{"last_run_at": &now}).Error + _ = s.archiveCompletedSubscription(ctx, sub, availability) + if queueResult.Queued > 0 { + s.hub.Publish("subscription", map[string]any{ + "id": sub.ID, + "name": sub.Name, + "queued": queueResult.Queued, + "keyword": keyword, + "resources": queueResult.Resources, + }) + s.notifySubscriptionHit(sub, queueResult.Queued, queueResult.Resources) + return queueResult.Queued, nil + } + if queueResult.LastEnqueueErr != nil { + return 0, fmt.Errorf("找到 PT 资源但加入下载器失败: %w", queueResult.LastEnqueueErr) + } + if s.log != nil { + fields := subscriptionSiteSearchLogFields(sub, keyword) + fields = appendSiteSearchSelectionLogFields(fields, selectionStats) + fields = appendAvailabilityLogFields(fields, availability) + fields = append(fields, zap.Int("queued", queueResult.Queued)) + s.log.Info("site-search subscription no candidate queued", fields...) + } + return 0, nil +} + +func (s *SubscriptionService) searchSubscriptionSites(ctx context.Context, sub *model.Subscription, keywords []string) ([]SearchResult, error) { var ( results []SearchResult lastSearchErr error @@ -55,20 +114,23 @@ func (s *SubscriptionService) runSiteSearch(ctx context.Context, sub *model.Subs } results = dedupeSiteSearchResults(results) if len(results) == 0 && lastSearchErr != nil && searchErrors == len(keywords) { - return 0, lastSearchErr + return nil, lastSearchErr } - if len(results) == 0 { - if s.log != nil { - fields := subscriptionSiteSearchLogFields(sub, keyword) - fields = append(fields, zap.Int("results_count", 0)) - s.log.Info("site-search subscription no results", fields...) - } - now := time.Now() - _ = s.repo.DB.Model(sub).Updates(map[string]any{"last_run_at": &now}).Error - return 0, nil - } - s.updateSubscriptionTotalEpisodes(ctx, sub, s.resolveSubscriptionTotalEpisodes(ctx, sub, inferSearchTotalEpisodes(results, sub))) + return results, nil +} +func (s *SubscriptionService) finishSiteSearchNoResults(sub *model.Subscription, keyword string) (int, error) { + if s.log != nil { + fields := subscriptionSiteSearchLogFields(sub, keyword) + fields = append(fields, zap.Int("results_count", 0)) + s.log.Info("site-search subscription no results", fields...) + } + now := time.Now() + _ = s.repo.DB.Model(sub).Updates(map[string]any{"last_run_at": &now}).Error + return 0, nil +} + +func (s *SubscriptionService) loadSiteSearchSeen(ctx context.Context, sub *model.Subscription) (string, []string, map[string]struct{}) { guidKey := fmt.Sprintf("subscription.%s.seen", sub.ID) seenRaw, _ := s.repo.Setting.Get(ctx, guidKey) seen := splitNonEmpty(seenRaw) @@ -76,171 +138,12 @@ func (s *SubscriptionService) runSiteSearch(ctx context.Context, sub *model.Subs for _, g := range seen { seenSet[g] = struct{}{} } - - availability := mergeLocalAvailability( - SubscriptionLocalAvailability(ctx, s.repo, sub), - s.pendingDownloadAvailability(ctx, sub), - ) - candidates, selectionStats := selectSiteSearchCandidatesWithStats(results, sub, seenSet, availability) - if s.log != nil { - fields := subscriptionSiteSearchLogFields(sub, keyword) - fields = appendSiteSearchSelectionLogFields(fields, selectionStats) - fields = appendAvailabilityLogFields(fields, availability) - s.log.Info("site-search subscription selection summary", fields...) - } - var lastEnqueueErr error - queued := 0 - var resources []string - for _, candidate := range candidates { - item := candidate.Item - matchText := subscriptionSearchResultText(item) - mediaType, mediaCategory := s.classifySubscriptionItem(ctx, sub, matchText, item.Category) - if s.shouldSkipExistingTorrent(ctx, mediaType, candidate) { - addSiteSearchCandidateAvailability(candidate, &availability) - seen = append(seen, candidate.GUID) - seenSet[candidate.GUID] = struct{}{} - if s.log != nil { - fields := subscriptionSiteSearchLogFields(sub, keyword) - fields = append(fields, - zap.String("reason", "existing_torrent"), - zap.String("title", item.Title), - zap.String("subtitle", item.Subtitle), - zap.String("site", firstNonEmpty(item.SiteName, item.SiteID)), - zap.String("site_category", item.Category), - zap.Int("season", candidate.Season), - zap.Int("episode", candidate.Episode), - zap.Bool("pack", candidate.Pack), - zap.String("media_type", mediaType), - ) - s.log.Info("site-search subscription candidate skipped", fields...) - } - continue - } - realURL := s.site.ResolveDownloadURL(ctx, candidate.Download) - savePath := s.resolveSubscriptionSavePath(ctx, sub, mediaType, mediaCategory) - if s.downloadPathHasCandidate(ctx, sub, matchText, savePath) { - addSiteSearchCandidateAvailability(candidate, &availability) - seen = append(seen, candidate.GUID) - seenSet[candidate.GUID] = struct{}{} - if s.log != nil { - fields := subscriptionSiteSearchLogFields(sub, keyword) - fields = append(fields, - zap.String("reason", "download_path_has_candidate"), - zap.String("title", item.Title), - zap.String("subtitle", item.Subtitle), - zap.String("site", firstNonEmpty(item.SiteName, item.SiteID)), - zap.String("site_category", item.Category), - zap.Int("season", candidate.Season), - zap.Int("episode", candidate.Episode), - zap.Bool("pack", candidate.Pack), - zap.String("media_type", mediaType), - zap.String("media_category", mediaCategory), - zap.String("save_path", savePath), - ) - s.log.Info("site-search subscription candidate skipped", fields...) - } - continue - } - if _, err := s.downloads.AddDownloadWithMeta(ctx, sub.UserID, realURL, savePath, DownloadTaskMeta{ - SubscriptionID: sub.ID, - Title: firstNonEmpty(item.Title, sub.Name), - PosterURL: sub.PosterURL, - BackdropURL: sub.BackdropURL, - Overview: sub.Overview, - MediaType: mediaType, - MediaCategory: mediaCategory, - SourceCategory: item.Category, - AllowExistingLibrary: sub.WashEnabled, - }); err != nil { - if IsDownloadDedupError(err) { - addSiteSearchCandidateAvailability(candidate, &availability) - seen = append(seen, candidate.GUID) - seenSet[candidate.GUID] = struct{}{} - if s.log != nil { - fields := subscriptionSiteSearchLogFields(sub, keyword) - fields = append(fields, - zap.String("reason", "download_dedup"), - zap.String("title", item.Title), - zap.String("subtitle", item.Subtitle), - zap.String("site", firstNonEmpty(item.SiteName, item.SiteID)), - zap.String("site_category", item.Category), - zap.Int("season", candidate.Season), - zap.Int("episode", candidate.Episode), - zap.Bool("pack", candidate.Pack), - zap.String("media_type", mediaType), - zap.String("media_category", mediaCategory), - zap.String("save_path", savePath), - ) - s.log.Info("site-search subscription candidate skipped", fields...) - } - continue - } - lastEnqueueErr = err - s.log.Warn("site-search subscription enqueue failed", - zap.String("subscription_id", sub.ID), - zap.String("subscription", sub.Name), - zap.String("keyword", keyword), - zap.String("title", item.Title), - zap.String("subtitle", item.Subtitle), - zap.String("site", firstNonEmpty(item.SiteName, item.SiteID)), - zap.String("site_category", item.Category), - zap.String("media_type", mediaType), - zap.String("media_category", mediaCategory), - zap.String("save_path", savePath), - zap.Error(err)) - continue - } - queued++ - addSiteSearchCandidateAvailability(candidate, &availability) - resources = append(resources, item.Title) - seen = append(seen, candidate.GUID) - seenSet[candidate.GUID] = struct{}{} - if s.log != nil { - fields := subscriptionSiteSearchLogFields(sub, keyword) - fields = append(fields, - zap.String("title", item.Title), - zap.String("subtitle", item.Subtitle), - zap.String("site", firstNonEmpty(item.SiteName, item.SiteID)), - zap.String("site_category", item.Category), - zap.Int("season", candidate.Season), - zap.Int("episode", candidate.Episode), - zap.Bool("pack", candidate.Pack), - zap.Int("score", candidate.Score), - zap.String("media_type", mediaType), - zap.String("media_category", mediaCategory), - zap.String("save_path", savePath), - ) - s.log.Info("site-search subscription candidate queued", fields...) - } - } - availability = s.finalizePendingAvailability(sub, availability) - if len(seen) > 200 { - seen = seen[len(seen)-200:] - } - _ = s.repo.Setting.Set(ctx, guidKey, strings.Join(seen, "\n")) - now := time.Now() - _ = s.repo.DB.Model(sub).Updates(map[string]any{"last_run_at": &now}).Error - _ = s.archiveCompletedSubscription(ctx, sub, availability) - if queued > 0 { - s.hub.Publish("subscription", map[string]any{ - "id": sub.ID, - "name": sub.Name, - "queued": queued, - "keyword": keyword, - "resources": resources, - }) - s.notifySubscriptionHit(sub, queued, resources) - return queued, nil - } - if lastEnqueueErr != nil { - return 0, fmt.Errorf("找到 PT 资源但加入下载器失败: %w", lastEnqueueErr) - } - if s.log != nil { - fields := subscriptionSiteSearchLogFields(sub, keyword) - fields = appendSiteSearchSelectionLogFields(fields, selectionStats) - fields = appendAvailabilityLogFields(fields, availability) - fields = append(fields, zap.Int("queued", queued)) - s.log.Info("site-search subscription no candidate queued", fields...) - } - return 0, nil + return guidKey, seen, seenSet +} + +func trimSiteSearchSeen(seen []string) []string { + if len(seen) <= 200 { + return seen + } + return seen[len(seen)-200:] } diff --git a/internal/service/subscription_site_search_enqueue.go b/internal/service/subscription_site_search_enqueue.go new file mode 100644 index 0000000..63519dc --- /dev/null +++ b/internal/service/subscription_site_search_enqueue.go @@ -0,0 +1,150 @@ +package service + +import ( + "context" + + "go.uber.org/zap" + + "github.com/ShukeBta/MediaStationGo/internal/model" +) + +type siteSearchRunState struct { + Keyword string + Seen []string + SeenSet map[string]struct{} + Availability LocalAvailability +} + +type siteSearchQueueResult struct { + Queued int + Resources []string + LastEnqueueErr error +} + +func (s *SubscriptionService) enqueueSiteSearchCandidates(ctx context.Context, sub *model.Subscription, candidates []siteSearchCandidate, state *siteSearchRunState) siteSearchQueueResult { + var result siteSearchQueueResult + for _, candidate := range candidates { + title, err := s.enqueueSiteSearchCandidate(ctx, sub, candidate, state) + if err != nil { + result.LastEnqueueErr = err + continue + } + if title == "" { + continue + } + result.Queued++ + result.Resources = append(result.Resources, title) + } + return result +} + +func (s *SubscriptionService) enqueueSiteSearchCandidate(ctx context.Context, sub *model.Subscription, candidate siteSearchCandidate, state *siteSearchRunState) (string, error) { + item := candidate.Item + matchText := subscriptionSearchResultText(item) + mediaType, mediaCategory := s.classifySubscriptionItem(ctx, sub, matchText, item.Category) + if s.shouldSkipExistingTorrent(ctx, mediaType, candidate) { + state.markCandidateAvailable(candidate) + s.logSiteSearchCandidateSkipped(sub, state, candidate, "existing_torrent", mediaType, "", "") + return "", nil + } + + realURL := s.site.ResolveDownloadURL(ctx, candidate.Download) + savePath := s.resolveSubscriptionSavePath(ctx, sub, mediaType, mediaCategory) + if s.downloadPathHasCandidate(ctx, sub, matchText, savePath) { + state.markCandidateAvailable(candidate) + s.logSiteSearchCandidateSkipped(sub, state, candidate, "download_path_has_candidate", mediaType, mediaCategory, savePath) + return "", nil + } + + if _, err := s.downloads.AddDownloadWithMeta(ctx, sub.UserID, realURL, savePath, DownloadTaskMeta{ + SubscriptionID: sub.ID, + Title: firstNonEmpty(item.Title, sub.Name), + PosterURL: sub.PosterURL, + BackdropURL: sub.BackdropURL, + Overview: sub.Overview, + MediaType: mediaType, + MediaCategory: mediaCategory, + SourceCategory: item.Category, + AllowExistingLibrary: sub.WashEnabled, + }); err != nil { + if IsDownloadDedupError(err) { + state.markCandidateAvailable(candidate) + s.logSiteSearchCandidateSkipped(sub, state, candidate, "download_dedup", mediaType, mediaCategory, savePath) + return "", nil + } + s.logSiteSearchEnqueueFailed(sub, state, candidate, mediaType, mediaCategory, savePath, err) + return "", err + } + + state.markCandidateAvailable(candidate) + s.logSiteSearchCandidateQueued(sub, state, candidate, mediaType, mediaCategory, savePath) + return item.Title, nil +} + +func (state *siteSearchRunState) markCandidateAvailable(candidate siteSearchCandidate) { + addSiteSearchCandidateAvailability(candidate, &state.Availability) + state.Seen = append(state.Seen, candidate.GUID) + if state.SeenSet != nil { + state.SeenSet[candidate.GUID] = struct{}{} + } +} + +func (s *SubscriptionService) logSiteSearchCandidateSkipped(sub *model.Subscription, state *siteSearchRunState, candidate siteSearchCandidate, reason, mediaType, mediaCategory, savePath string) { + if s.log == nil { + return + } + fields := subscriptionSiteSearchLogFields(sub, state.Keyword) + fields = append(fields, zap.String("reason", reason)) + fields = appendSiteSearchCandidateLogFields(fields, candidate) + fields = append(fields, zap.String("media_type", mediaType)) + if mediaCategory != "" { + fields = append(fields, zap.String("media_category", mediaCategory)) + } + if savePath != "" { + fields = append(fields, zap.String("save_path", savePath)) + } + s.log.Info("site-search subscription candidate skipped", fields...) +} + +func (s *SubscriptionService) logSiteSearchCandidateQueued(sub *model.Subscription, state *siteSearchRunState, candidate siteSearchCandidate, mediaType, mediaCategory, savePath string) { + if s.log == nil { + return + } + fields := subscriptionSiteSearchLogFields(sub, state.Keyword) + fields = appendSiteSearchCandidateLogFields(fields, candidate) + fields = append(fields, + zap.Int("score", candidate.Score), + zap.String("media_type", mediaType), + zap.String("media_category", mediaCategory), + zap.String("save_path", savePath), + ) + s.log.Info("site-search subscription candidate queued", fields...) +} + +func (s *SubscriptionService) logSiteSearchEnqueueFailed(sub *model.Subscription, state *siteSearchRunState, candidate siteSearchCandidate, mediaType, mediaCategory, savePath string, err error) { + if s.log == nil { + return + } + fields := subscriptionSiteSearchLogFields(sub, state.Keyword) + fields = appendSiteSearchCandidateLogFields(fields, candidate) + fields = append(fields, + zap.String("media_type", mediaType), + zap.String("media_category", mediaCategory), + zap.String("save_path", savePath), + zap.Error(err), + ) + s.log.Warn("site-search subscription enqueue failed", fields...) +} + +func appendSiteSearchCandidateLogFields(fields []zap.Field, candidate siteSearchCandidate) []zap.Field { + item := candidate.Item + return append(fields, + zap.String("title", item.Title), + zap.String("subtitle", item.Subtitle), + zap.String("site", firstNonEmpty(item.SiteName, item.SiteID)), + zap.String("site_category", item.Category), + zap.Int("season", candidate.Season), + zap.Int("episode", candidate.Episode), + zap.Bool("pack", candidate.Pack), + ) +}